0


Flink将数据流导入Doris

1.Maven依赖导入

<dependency>
  <groupId>org.apache.doris</groupId>
  <artifactId>flink-doris-connector-1.16</artifactId>
  <version>1.4.0</version>
</dependency> 

1.1版本兼容

Connector VersionFlink VersionDoris VersionJava VersionScala Version1.0.31.11+0.15+82.11,2.121.1.11.141.0+82.11,2.121.2.11.151.0+8-1.3.01.161.0+8-1.4.01.15,1.16,1.171.0+8-

2.代码实现(DataStream)

    DorisSink是通过StreamLoad向Doris写入数据,DataStream写入时,支持不同的序列化方法。这里以String 数据流为例。

2.1Doris表准备

CREATE TABLE IF NOT EXISTS demo.test
(
    `user_id` LARGEINT NOT NULL COMMENT "用户id",
    `cost` BIGINT SUM DEFAULT "0" COMMENT "用户总消费"
)
AGGREGATE KEY(`user_id`)
DISTRIBUTED BY HASH(`user_id`) BUCKETS 1
PROPERTIES (
    "replication_allocation" = "tag.location.default: 1"
);

2.2 java实现

   @Test
    public void test01() throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);
        env.enableCheckpointing(10000);
        Properties properties = new Properties();
        DorisSink<String> dorisSink = DorisSink.<String>builder()  // 创建一个DorisSink对象用于将数据写入Doris数据库
                .setDorisOptions(  // 设置Doris数据库的连接选项
                        DorisOptions.builder()  // 创建一个DorisOptions对象
                                .setFenodes("FE_IP:HTTP_PORT")  // 设置Fenode服务器地址
                                .setTableIdentifier("db.table")  // 设置表的标识符
                                .setUsername("root")  // 设置用户名
                                .setPassword("123456")  // 设置密码
                                .build()  // 构建DorisOptions对象
                )
                .setDorisExecutionOptions(  // 设置Doris数据库的执行选项
                        DorisExecutionOptions.builder()  // 创建一个DorisExecutionOptions对象
                                .setLabelPrefix("doris")  // 设置标签前缀
                                .setDeletable(false)  // 设置是否可删除
                                .setStreamLoadProp(properties)  // 设置流加载属性
                                .build()  // 构建DorisExecutionOptions对象
                )
                .setDorisReadOptions(  // 设置Doris数据库的读取选项
                        DorisReadOptions.builder().build()  // 创建一个DorisReadOptions对象并构建
                ).setSerializer(new SimpleStringSerializer())  // 设置数据序列化方式为SimpleStringSerializer
                .build();  // 构建DorisSink对象

        DataStreamSource<String> streamSource = env.fromElements("50" + "\t" + 1);

        streamSource.print();
        streamSource.sinkTo(dorisSink);
        env.execute();
    }

2.3 配置项

  • setLabelPrefix:Stream load导入使用的label前缀。2pc场景下要求全局唯一 ,用来保证Flink的EOS语义。
  • setStreamLoadProp:

常用配置:可以在properties中添加

  • 自定义字符串列分隔符:sink.properties.column_separator' = ', '
  • 特殊字符作为分隔符:'sink.properties.escape_delimiters' = 'true'

如果是JSON格式导入需要加2个配置项

  • 'sink.properties.format' = 'json'
  • 'sink.properties.read_json_by_line' = 'true'
  • setDeletable:是否开启批量删除功能。只支持 Unique 模型(默认为true)
标签: flink 大数据

本文转载自: https://blog.csdn.net/WB231444/article/details/135861371
版权归原作者 星辰境末 所有, 如有侵权,请联系我们删除。

“Flink将数据流导入Doris”的评论:

还没有评论