资讯动态

Doris与Kafka/Flink实时集成:Routine Load、Stream Load与实时数仓链路

发布时间:2026/9/28 20:47:08 来源:尧图企业网站定制
引言Apache Doris原名Palo是一个现代化的MPP分析型数据库具备高并发、低延迟的查询性能。在现代数据架构中实时数据集成是构建高效数据仓库的关键环节。Doris通过与Kafka和Flink的集成能够实现实时数据的无缝接入与处理为企业提供实时的数据分析能力。本文将详细介绍Doris通过Routine Load和Stream Load两种方式与Kafka集成的实现方法并结合Flink构建完整的实时数仓链路。Routine Load 详细介绍与实现Routine Load是Doris提供的一种持续从消息队列如Kafka拉取数据并导入的方式特别适合处理持续不断的数据流。2.1 工作原理Routine Load通过创建一个异步任务定期从Kafka中读取数据并导入到Doris表中。该任务会持续运行直到被手动停止。Doris使用一个名为broker的组件连接到Kafka集群获取数据并导入。2.2 适用场景Routine Load适用于以下场景持续不断的数据流导入数据量较大但实时性要求不是最高的场景能够容忍少量数据延迟的场景无需频繁启动/停止导入任务的情况2.3 实现步骤实现Routine Load的基本步骤如下在Doris中创建导入任务CREATE ROUTINE LOAD example_db.example_load ON TABLE example_table COLUMNS(k1, k2, k3), PROPERTIES ( format json, strip_outer_array true ) FROM KAFKA ( kafka_broker kafka1:9092,kafka2:9092, kafka_topic example_topic, property.group.id example_group );查看任务状态SHOW ROUTINE LOAD WHERE NAME example_load;暂停任务PAUSE ROUTINE LOAD FOR example_db.example_load;恢复任务RESUME ROUTINE LOAD FOR example_db.example_load;停止任务STOP ROUTINE LOAD FOR example_db.example_load;2.4 优缺点分析特点描述优点1. 持续运行无需手动干预br2. 支持高吞吐量数据导入br3. 自动重试机制保证数据可靠性br4. 支持数据过滤和转换缺点1. 启动和停止操作较为复杂br2. 无法精确控制导入时间点br3. 任务状态管理较为繁琐br4. 配置参数较多调试成本高Stream Load 详细介绍与实现Stream Load是Doris提供的一种通过HTTP协议导入数据的方式适合需要手动触发或通过程序控制的数据导入场景。3.1 工作原理Stream Load通过发送HTTP请求将数据文件导入到Doris表中。客户端可以构造包含数据内容的HTTP请求Doris收到请求后解析数据并导入。这种方式不需要额外的服务组件直接通过Doris的HTTP服务完成导入。3.2 适用场景Stream Load适用于以下场景需要精确控制导入时间的场景数据量相对较小但要求实时性高的场景程序化控制导入过程的场景需要快速验证数据正确性的场景3.3 实现步骤实现Stream Load的基本步骤如下使用curl命令导入数据curl --location-trusted -u user:password -H label:stream_load_20230601 -H Content-Type: text/plain -T data.txt http://doris_host:8030/api/example_db/example_table/_stream_load参数说明user:passwordDoris认证信息label导入任务标签用于唯一标识一次导入Content-Type数据格式支持多种格式data.txt包含数据的文件http://doris_host:8030/api/example_db/example_table/_stream_loadStream Load API地址程序化实现Python示例import requests def stream_load_to_doris(data, doris_url, user, password, table): headers { label: stream_load_ str(int(time.time())), Content-Type: text/plain } response requests.post( fhttp://{doris_url}/api/{table}, auth(user, password), headersheaders, datadata ) return response.json() # 使用示例 data 1,example1,2023-06-01\n2,example2,2023-06-01 result stream_load_to_doris(data, doris_host:8030, user, password, example_db.example_table) print(result)3.4 优缺点分析特点描述优点1. 实现简单无需额外组件br2. 可精确控制导入时机br3. 支持各种编程语言调用br4. 导入状态反馈及时缺点1. 需要客户端构造数据br2. 不适合持续不断的数据流br3. 大数据量导入性能较差br4. 需要处理网络异常情况实时数仓链路构建结合Routine Load/Stream Load与Flink可以构建完整的实时数仓链路。下面通过mermaid流程图展示整体架构StreamLoad方式RoutineLoad方式Flink处理数据接入实时计算结果输出数据源KafkaFlink处理数据清洗/转换Doris存储实时查询Doris Routine Load任务客户端程序Stream Load API构建实时数仓链路的步骤如下数据接入层使用Kafka作为消息队列接收来自业务系统的实时数据实时处理层通过Flink从Kafka读取数据进行清洗、转换和聚合处理数据存储层处理后的数据通过Routine Load或Stream Load导入Doris查询服务层通过Doris提供高性能的实时查询服务一个完整的Flink处理示例import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer; import org.apache.flink.streaming.util.serialization.SimpleStringSchema; public class DorisRealtimeETL { public static void main(String[] args) throws Exception { // 创建执行环境 final StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); // 配置Kafka消费者 properties.setProperty(bootstrap.servers, kafka1:9092,kafka2:9092); properties.setProperty(group.id, doris_etl); FlinkKafkaConsumerString kafkaSource new FlinkKafkaConsumer( example_topic, new SimpleStringSchema(), properties ); // 添加数据源 DataStreamString stream env.addSource(kafkaSource); // 数据处理逻辑示例简单解析和转换 DataStreamString processedStream stream.map(value - { // 解析JSON数据并进行转换 // 这里只是示例实际应根据业务需求编写处理逻辑 return processAndTransform(value); }); // 输出到Doris通过自定义Sink或调用Stream Load API processedStream.addSink(new DorisSink()); // 执行任务 env.execute(Doris Realtime ETL); } private static String processAndTransform(String input) { // 数据处理逻辑 return transformedData; } }完整示例与注意事项5.1 最小示例以下是一个完整的Doris与Kafka/Flink集成的最小示例创建Doris表CREATE TABLE example_db.user_events ( user_id BIGINT, event_type VARCHAR(32), event_time DATETIME, event_data VARCHAR(1024) ) ENGINEOLAP PRIMARY KEY(user_id, event_time) DISTRIBUTED BY HASH(user_id) BUCKETS 10;创建Routine Load任务CREATE ROUTINE LOAD example_db.user_events_load ON TABLE example_db.user_events COLUMNS(user_id, event_type, event_time, event_data), PROPERTIES ( format json, jsonpaths $.user_id, $.event_type, $.event_time, $.event_data ) FROM KAFKA ( kafka_broker kafka1:9092,kafka2:9092, kafka_topic user_events, property.group.id user_events_group );Flink处理代码简化版import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer; import org.apache.flink.streaming.util.serialization.JSONKeyValueDeserializationSchema; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.JsonNode; public class SimpleFlinkJob { public static void main(String[] args) throws Exception { final StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); Properties properties new Properties(); properties.setProperty(bootstrap.servers, kafka1:9092,kafka2:9092); properties.setProperty(group.id, user_events_group); FlinkKafkaConsumerJsonNode kafkaSource new FlinkKafkaConsumer( user_events, new JSONKeyValueDeserializationSchema(true), properties ); env.addSource(kafkaSource) .addSink(new DorisSink(example_db.user_events)); env.execute(Simple Flink to Doris); } }Doris Sink实现简化版import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer; import org.apache.flink.streaming.connectors.kafka.KafkaSerializationSchema; import org.apache.kafka.clients.producer.ProducerRecord; public class DorisSink extends RichSinkFunctionJsonNode { private final String table; private transient HttpPost httpPost; public DorisSink(String table) { this.table table; } Override public void open(Configuration parameters) throws Exception { // 初始化HTTP请求 httpPost new HttpPost(http://doris_host:8030/api/ table /_stream_load); httpPost.setHeader(Content-Type, application/json); httpPost.setHeader(Authorization, Basic Base64.encodeBase64String(user:password.getBytes())); } Override public void invoke(JsonNode value, Context context) throws Exception { // 发送数据到Doris String requestBody value.toString(); httpPost.setEntity(new StringEntity(requestBody)); try (CloseableHttpResponse response httpClient.execute(httpPost)) { // 处理响应 EntityUtils.consume(response.getEntity()); } } Override public void close() throws Exception { if (httpPost ! null) { httpPost.releaseConnection(); } } }5.2 注意事项性能优化Doris表设计中合理设置分桶数量针对查询模式优化列存储批量导入时调整超时参数数据一致性使用唯一的label确保数据不重复设置合适的重试机制监控导入失败情况资源管理合理设置Flink并行度控制Doris导入任务数量监控系统资源使用情况错误处理实现完善的错误捕获机制设置合理的重试策略记录详细的日志信息安全考虑使用HTTPS连接Doris定期更换认证凭据实施最小权限原则通过以上内容我们了解了Doris与Kafka/Flink实时集成的实现方法包括Routine Load和Stream Load两种数据导入方式以及如何构建完整的实时数仓链路。结合提供的示例代码和注意事项开发者可以根据实际业务需求快速搭建实时数据集成平台。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑