
StarRocks Routine Load 实战从 Kafka 持续加载 CSV、JSON 与 Avro 数据【免费下载链接】starrocksThe worlds fastest open query engine for sub-second analytics both on and off the data lakehouse. With the flexibility to support nearly any scenario, StarRocks provides best-in-class performance for multi-dimensional analytics, real-time analytics, and ad-hoc queries. A Linux Foundation project.项目地址: https://gitcode.com/GitHub_Trending/st/starrocks本文基于 StarRocks 仓库 RoutineLoad.md 文档撰写并结合fe/fe-core/src/main/java/com/starrocks/load/routineload/下的调度器与任务源码展开。Routine Load 是 StarRocks 内置的常驻型long-running流式导入方案只需一条CREATE ROUTINE LOAD语句即可让 FE 持续把 Kafka 主题topic中的消息拆分成一批批并发加载任务load task写入 StarRocks并保证 exactly-once 语义。读完本文你将掌握 Routine Load 的核心概念与完整工作流、CSV/JSON/Avro 三种数据格式的建表与建任务实操、源消息元数据topic/partition/offset/key/headers的注入写法以及任务从暂停、恢复、修改到停止的完整生命周期管理。Routine Load 是什么Routine Load 是 StarRocks 提供的、面向 Kafka 消息流的持续数据加载方案。与一次性执行的 Stream Load 不同Routine Load 任务一旦创建便会常驻运行只要任务状态为RUNNING它就会持续生成一批又一批的加载任务消费 Kafka 主题中全部或部分分区partition的消息并将其加载进 StarRocks 表。Routine Load 的典型适用场景包括实时日志、埋点事件、订单流水等持续到达的 Kafka 消息流需要以秒级延迟将流数据刷入 OLAP 表做多维分析与即席查询需要在加载过程中完成字段映射、类型转换如时间戳转 DATE等 ETL 操作。Routine Load 具备以下核心能力Exactly-once 语义保证加载进 StarRocks 的数据既不丢失也不重复源码中 FE 会根据任务结果重试失败任务详见 RoutineLoadJob.java加载时数据转换支持在加载过程中对数据进行 ETL 变换详见 Transform data at loading加载时数据变更支持 UPSERT、DELETE 等写入操作详见 Change data through loading多格式支持支持 CSV、JSON 以及v3.0.1 起Avro 格式数据。权限要求执行 Routine Load 的用户必须对目标表拥有INSERT权限。若无权限需按 GRANT 文档执行GRANT INSERT ON TABLE table_name IN DATABASE database_name TO { ROLE role_name | USER user_identity}授权。更多权限用法参见 account-management。核心概念与工作流程术语Load Job 与 Load TaskLoad job加载作业Routine Load 作业是常驻型任务。只要其状态为RUNNING作业就会持续生成一个或多个并发的加载任务消费 Kafka 集群某个主题中的消息并加载进 StarRocks。一个表上可以同时存在多个加载作业。Load task加载任务加载作业会按一定规则被拆分成多个加载任务。加载任务是数据加载的基本单元每个任务是一个独立事件其加载机制基于 Stream Load 实现。多个加载任务并发地从主题的不同分区消费消息并写入 StarRocks。下图展示了 Routine Load 的整体架构FE 管理 Routine Load Job 的元数据与生命周期将其拆分为多个 task 并下发给 BE 节点执行实际的数据写入。完整工作流程Routine Load 的运行分为以下四个阶段创建 Routine Load 作业通过 CREATE ROUTINE LOAD 语句创建作业FE 解析该语句并按照所指定的属性创建作业。从源码看RoutineLoadMgr.createRoutineLoadJob()负责作业注册与状态初始化RoutineLoadMgr.java。FE 将作业拆分为多个加载任务FE 按规则拆分任务每个加载任务是一个独立事务。拆分规则如下FE 依据期望并发数desired_concurrent_number、Kafka 主题分区数以及存活 BE 节点数计算实际并发数FE 按计算出的实际并发数拆分作业并将任务排入任务队列。Kafka 主题由多个分区组成主题分区与加载任务的对应关系为一个分区唯一分配给一个加载任务该分区的所有消息由该任务消费一个加载任务可以消费一个或多个分区的消息所有分区在加载任务之间均匀分布。多个加载任务并发消费多个 Kafka 主题分区的消息并加载进 StarRocksFE 调度并提交加载任务FE 定时调度队列中的任务并将其分配给选定的 Coordinator BE 节点。任务调度间隔由配置项max_batch_interval定义FE 将任务均匀分发给所有 BE 节点。Coordinator BE 启动加载任务消费分区中的消息、解析并过滤数据。一个加载任务会持续运行直到消费的消息量达到预定义上限或达到预定义时间上限。消息批量大小与时间上限分别由 FE 配置项max_routine_load_batch_size与routine_load_task_consume_second定义详见 FE Configuration。随后 Coordinator BE 将消息分发给 Executor BE 节点Executor BE 将消息写入磁盘。安全协议StarRocks 支持通过 SASL_SSL、SASL_PLAINTEXT、SSL 与 PLAINTEXT 等安全协议访问 Kafka。本文示例均以 PLAINTEXT 连接为例如需使用其他安全协议参见 CREATE ROUTINE LOAD。FE 持续生成新的加载任务Executor BE 将数据写入磁盘后Coordinator BE 向 FE 上报该加载任务的执行结果。FE 据此生成新的加载任务持续加载数据或重试失败的任务从而保证加载进 StarRocks 的数据既不丢失也不重复。从源码看调度实现在 FE 侧Routine Load 的调度逻辑分布在fe/fe-core/src/main/java/com/starrocks/load/routineload/目录下RoutineLoadScheduler继承LeaderDaemon的常驻守护线程周期性地取出状态为NEED_SCHEDULE的作业进行调度并驱动NEED_SCHEDULE - RUNNING的状态迁移RoutineLoadScheduler.javaRoutineLoadTaskScheduler负责任务队列的调度与提交RoutineLoadTaskScheduler.javaKafkaRoutineLoadJob.calculateCurrentConcurrentTaskNum()直接体现了取最小值的并发计算逻辑即实际并发数 min(分区数, desired_concurrent_number, 存活 BE 数, max_routine_load_task_concurrent_num)KafkaRoutineLoadJob.java。在 BE 侧消息消费由be/src/data_workflows/load/routine_load/下的data_consumer数据消费者与kafka_consumer_pipeKafka 消费管道实现负责从 Kafka 拉取消息并送入 Stream Load 管道写入磁盘。支持的数据格式Routine Load 目前支持从 Kafka 集群消费 CSV、JSON 以及 Avrov3.0.1 起格式的数据。CSV 数据注意事项文本分隔符可以使用不超过 50 字节的 UTF-8 字符串例如逗号,、制表符\t或竖线|空值NULL使用\N表示。例如数据文件包含三列某条记录第一列和第三列有数据、第二列无数据则该记录必须写成a,\N,b而非a,,b——a,,b表示第二列是空字符串。创建 Routine Load 作业三种格式实战下面通过三个完整示例分别演示消费 Kafka 中的 CSV、JSON 与 Avro 格式数据并加载进 StarRocks。完整语法与参数说明见 CREATE ROUTINE LOAD。加载 CSV 格式数据准备数据集假设 Kafka 集群的ordertest1主题中存在 CSV 格式数据集每条消息包含六个字段订单 ID、支付日期、客户姓名、国籍、性别、价格。2020050802,2020-05-08,Johann Georg Faust,Deutschland,male,895 2020050802,2020-05-08,Julien Sorel,France,male,893 2020050803,2020-05-08,Dorian Grey,UK,male,1262 2020050901,2020-05-09,Anna Karenina,Russia,female,175 2020051001,2020-05-10,Tess Durbeyfield,US,female,986 2020051101,2020-05-11,Edogawa Conan,japan,male,8924创建表根据 CSV 数据的字段在数据库example_db中创建表example_tbl1。以下示例创建了 5 个字段的表不包含CSV 数据中的客户性别字段。CREATE TABLE example_db.example_tbl1 ( order_id bigint NOT NULL COMMENT Order ID, pay_dt date NOT NULL COMMENT Payment date, customer_name varchar(26) NULL COMMENT Customer name, nationality varchar(26) NULL COMMENT Nationality, pricedouble NULL COMMENT Price ) ENGINEOLAP DUPLICATE KEY (order_id,pay_dt) DISTRIBUTED BY HASH(order_id);说明自 v2.5.7 起StarRocks 在创建表或添加分区时可自动设置分桶数BUCKETS无需再手动指定。详见 set the number of buckets。提交 Routine Load 作业执行以下语句提交名为example_tbl1_ordertest1的 Routine Load 作业消费主题ordertest1中的消息并加载进表example_tbl1。加载任务从指定分区的初始 offset 开始消费消息。CREATE ROUTINE LOAD example_db.example_tbl1_ordertest1 ON example_tbl1 COLUMNS TERMINATED BY ,, COLUMNS (order_id, pay_dt, customer_name, nationality, temp_gender, price) PROPERTIES ( desired_concurrent_number 5 ) FROM KAFKA ( kafka_broker_list kafka_broker1_ip:kafka_broker1_port,kafka_broker2_ip:kafka_broker2_port, kafka_topic ordertest1, kafka_partitions 0,1,2,3,4, property.kafka_default_offsets OFFSET_BEGINNING );提交作业后可以执行 SHOW ROUTINE LOAD 语句查看作业状态。关键参数解读作业命名一张表上可能存在多个加载作业建议以对应的 Kafka 主题 提交时间命名作业便于区分每张表上的不同作业。列分隔符COLUMNS TERMINATED BY属性定义 CSV 数据的列分隔符默认值为\t。Kafka 主题分区与 offset可以通过kafka_partitions与kafka_offsets指定要消费的分区及起始 offset。例如让作业从主题ordertest1的分区0,1,2,3,4以各自不同的起始 offset 消费消息可配置如下kafka_partitions 0,1,2,3,4, kafka_offsets OFFSET_BEGINNING, OFFSET_END, 1000, 2000, 3000也可以使用property.kafka_default_offsets统一设置所有分区的默认 offsetkafka_partitions 0,1,2,3,4, property.kafka_default_offsets OFFSET_BEGINNING源码补充在 CreateRoutineLoadStmt.java 中OFFSET_BEGINNING被解析为-2、OFFSET_END被解析为-1内部以KafkaProgress.OFFSET_BEGINNING_VAL/OFFSET_END_VAL表示也支持直接填写数值型 offset见 KafkaProgress.java。另外kafka_broker_list与kafka_topic为必填属性源码中缺失时会直接抛出AnalysisException。数据映射与转换CSV 数据与 StarRocks 表之间的映射与转换关系通过COLUMNS参数指定。数据映射StarRocks 将 CSV 数据中的列按顺序提取并映射到COLUMNS参数声明的字段上再将COLUMNS中声明的字段按名称映射到 StarRocks 表的列上。数据转换由于示例不加载客户性别字段COLUMNS中的temp_gender作为该字段的占位符其余字段直接映射到example_tbl1的列。更多转换方式见 Transform data at loading。注意如果 CSV 数据中列的名称、数量与顺序和 StarRocks 表完全对应则无需指定COLUMNS参数。任务并发度当 Kafka 主题分区较多且 BE 节点充足时可以通过提高任务并发度加速加载。提高实际加载任务并发度的方法有二一是在创建作业时增大期望并发数desired_concurrent_number二是将 FE 动态配置项max_routine_load_task_concurrent_num单个 Routine Load 作业的最大并发任务数默认 5源码见 Config.java调大。实际任务并发数由存活 BE 节点数、指定的 Kafka 主题分区数、desired_concurrent_number与max_routine_load_task_concurrent_num四者的最小值决定。这一逻辑与源码中KafkaRoutineLoadJob.calculateCurrentConcurrentTaskNum()的实现完全一致KafkaRoutineLoadJob.java。本例中存活 BE 数为5、指定 Kafka 分区数为5、max_routine_load_task_concurrent_num为5因此将desired_concurrent_number从默认值3提升到5即可提高实际并发数。加载 JSON 格式数据准备数据集假设 Kafka 集群的ordertest2主题中存在 JSON 格式数据集包含六个 key商品 ID、客户姓名、国籍、支付时间与价格。另外需要将支付时间列转换为 DATE 类型并加载进 StarRocks 表的pay_dt列。{commodity_id: 1, customer_name: Mark Twain, country: US,pay_time: 1589191487,price: 875} {commodity_id: 2, customer_name: Oscar Wilde, country: UK,pay_time: 1589191487,price: 895} {commodity_id: 3, customer_name: Antoine de Saint-Exupéry,country: France,pay_time: 1589191487,price: 895}注意每行的一个 JSON 对象必须完整地放在一条 Kafka 消息中否则会返回 JSON 解析错误。创建表根据 JSON 数据的 key在数据库example_db中创建表example_tbl2。CREATE TABLE example_tbl2 ( commodity_id varchar(26) NULL COMMENT Commodity ID, customer_name varchar(26) NULL COMMENT Customer name, country varchar(26) NULL COMMENT Country, pay_time bigint(20) NULL COMMENT Payment time, pay_dt date NULL COMMENT Payment date, pricedouble SUM NULL COMMENT Price ) ENGINEOLAP AGGREGATE KEY(commodity_id,customer_name,country,pay_time,pay_dt) DISTRIBUTED BY HASH(commodity_id);提交 Routine Load 作业CREATE ROUTINE LOAD example_db.example_tbl2_ordertest2 ON example_tbl2 COLUMNS(commodity_id, customer_name, country, pay_time, price, pay_dtfrom_unixtime(pay_time, %Y%m%d)) PROPERTIES ( desired_concurrent_number 5, format json, jsonpaths [\$.commodity_id\,\$.customer_name\,\$.country\,\$.pay_time\,\$.price\] ) FROM KAFKA ( kafka_broker_list kafka_broker1_ip:kafka_broker1_port,kafka_broker2_ip:kafka_broker2_port, kafka_topic ordertest2, kafka_partitions 0,1,2,3,4, property.kafka_default_offsets OFFSET_BEGINNING );关键参数解读数据格式必须在PROPERTIES子句中指定format json以声明数据格式为 JSON。数据映射与转换JSON 数据与 StarRocks 表的映射与转换关系通过COLUMNS参数与jsonpaths属性共同指定。COLUMNS中字段的顺序必须与 JSON 数据中字段的顺序一致字段名称必须与 StarRocks 表的列一致jsonpaths用于从 JSON 数据中提取所需字段这些字段随后由COLUMNS命名。数据映射StarRocks 提取 JSON 数据中的 key 并映射到jsonpaths属性声明的 key 上将jsonpaths中声明的 key按顺序映射到COLUMNS参数声明的字段上再将COLUMNS中声明的字段按名称映射到 StarRocks 表的列上。数据转换由于示例需要将pay_time转换为 DATE 类型并加载进pay_dt列需要在COLUMNS中使用from_unixtime(pay_time, %Y%m%d)函数其余字段直接映射到example_tbl2的列。注意如果 JSON 对象中 key 的名称与数量完全匹配 StarRocks 表的字段则无需指定COLUMNS参数。加载 Avro 格式数据v3.0.1 起准备数据集Avro schema创建如下 Avro schema 文件avro_schema.avsc{ type: record, name: sensor_log, fields : [ {name: id, type: long}, {name: name, type: string}, {name: checked, type : boolean}, {name: data, type: double}, {name: sensor_type, type: {type: enum, name: sensor_type_enum, symbols : [TEMPERATURE, HUMIDITY, AIR-PRESSURE]}} ] }将 Avro schema 注册到 Schema Registry 中。Avro 数据准备 Avro 数据并发送到 Kafka 主题topic_0。创建表根据 Avro 数据的字段在目标数据库example_db中创建表sensor_log。表的列名必须与 Avro 数据中的字段名匹配类型映射关系见下文数据类型映射。CREATE TABLE example_db.sensor_log ( id bigint NOT NULL COMMENT sensor id, name varchar(26) NOT NULL COMMENT sensor name, checked boolean NOT NULL COMMENT checked, data double NULL COMMENT sensor data, sensor_type varchar(26) NOT NULL COMMENT sensor type ) ENGINEOLAP DUPLICATE KEY (id) DISTRIBUTED BY HASH(id);提交 Routine Load 作业CREATE ROUTINE LOAD example_db.sensor_log_load_job ON sensor_log PROPERTIES ( format avro ) FROM KAFKA ( kafka_broker_list kafka_broker1_ip:kafka_broker1_port,kafka_broker2_ip:kafka_broker2_port,..., confluent.schema.registry.url http://172.xx.xxx.xxx:8081, kafka_topic topic_0, kafka_partitions 0,1,2,3,4,5, property.kafka_default_offsets OFFSET_BEGINNING );关键参数解读数据格式必须在PROPERTIES子句中指定format avro以声明数据格式为 Avro。Schema Registry需要配置confluent.schema.registry.url指定注册了 Avro schema 的 Schema Registry 地址StarRocks 通过该 URL 获取 Avro schema。格式如下confluent.schema.registry.url http[s]://[schema-registry-api-key:schema-registry-api-secret]hostname|ip address[:port]该配置项在 FE 侧由KafkaRoutineLoadJob解析并存入数据源属性源码见 KafkaRoutineLoadJob.java。数据映射与转换Avro 数据与 StarRocks 表的映射与转换关系通过COLUMNS参数与jsonpaths属性共同指定。COLUMNS中字段的顺序必须与jsonpaths中字段的顺序一致字段名称必须与 StarRocks 表的列一致。jsonpaths用于从 Avro 数据中提取所需字段这些字段随后由COLUMNS命名。注意如果 Avro record 中字段的名称与数量和 StarRocks 表的列完全匹配则无需指定COLUMNS参数。数据类型映射Avro 数据字段与 StarRocks 表列之间的数据类型映射如下基本类型AvroStarRocksnulNULLbooleanBOOLEANintINTlongBIGINTfloatFLOATdoubleDOUBLEbytesSTRINGstringSTRING复杂类型AvroStarRocksrecordSTRUCT或将整个 RECORD 或其子字段以 JSON 形式加载enumsSTRINGarraysARRAYmapsMAP 或 JSONunion(T, null)NULLABLE(T)fixedSTRING限制目前 StarRocks 不支持 Avro schema 演进schema evolution每条 Kafka 消息必须只包含一条 Avro 数据记录。访问源消息元数据INCLUDE METADATA在加载 JSON 或 Avro 格式数据时除了消息载荷payload还可以从 Kafka/Pulsar 消息的元数据中填充目标列——包括主题topic、分区partition、偏移量offset、时间戳、key 和 headers。通过INCLUDE METADATA (...)子句将每个元数据 key 绑定到一个别名alias别名是一个普通的源列可以在COLUMNS中引用。该特性适用于审计记录某行来自哪个 topic/partition/offset、事件时间处理使用消息时间戳以及基于 header 值的路由等场景。语法INCLUDE METADATA ( metadata_key [AS alias] [, metadata_key [AS alias] ...] )INCLUDE METADATA是一个加载属性需要放在其他加载属性如COLUMNS、WHERE之间顺序任意且位于PROPERTIES和FROM子句之前。AS alias可省略省略时别名默认为metadata_key。别名在子句内必须唯一且不能与载荷字段、目标表列或保留列名冲突。源码补充该子句在 FE 侧由ImportMetadataStmtImportMetadataStmt.java表示每个元数据 key 绑定一个隐藏源列key 会在分析阶段按数据源校验并映射为对应的元数据类型。toSql()方法统一用于作业持久化与SHOW CREATE ROUTINE LOAD渲染保证持久化形态与展示形态一致。支持的元数据 key支持哪些 key 取决于数据源数据源Key类型说明KAFKATOPICVARCHAR主题名称。KAFKAPARTITIONINT分区号。KAFKAOFFSETBIGINT消息在分区内的偏移量。KAFKATIMESTAMP_MSBIGINT自 epoch 起的记录时间戳毫秒。broker 未上报时间戳时为NULL。KAFKAKEYVARCHAR消息 key原始字节。消息无 key 时为NULL。KAFKAHEADERSMAPVARCHAR, VARCHAR所有 headers 组成的 map。key 重复时最后一个值生效。PULSARTOPICVARCHAR作业消费的逻辑主题分区主题的-partition-N后缀不包含在内。PULSARPARTITIONINT分区索引从每条消息的主题名解析。非分区主题为NULL。PULSARKEYVARCHAR分区键。消息无 key 时为NULL。PULSARMESSAGE_IDVARCHAR消息 ID。PULSARPUBLISH_TIME_MSBIGINT自 epoch 起的发布时间毫秒。PULSAREVENT_TIME_MSBIGINT自 epoch 起的事件时间毫秒。生产者未设置时为NULL。PULSARPROPERTIESMAPVARCHAR, VARCHAR所有 properties 组成的 map。key 重复时最后一个值生效。如需读取单个 header/property 值可对HEADERS/PROPERTIESmap 使用element_at(headers_alias, name)key 重复时最后一个值生效key 不存在时返回NULL。header/property 值作为原始字节原样放入 VARCHAR不做 UTF-8 校验。使用注意事项INCLUDE METADATA适用于format jsonKafka 与 Pulsar和format avro仅 KafkaPulsar Routine Load 不支持 Avro。不适用于 CSV因为一条 CSV 消息可能展开成多行每条消息的元数据会变得有歧义。元数据别名是普通的源列可以在COLUMNS表达式中任意引用。如果载荷字段与元数据 key 同名请用AS alias指定不同的别名以避免歧义。OFFSET仅 Kafka 支持MESSAGE_ID仅 Pulsar 支持使用数据源不支持的 key 会报错并列出该数据源支持的 key。HEADERS/PROPERTIES是MAP\VARCHAR, VARCHAR。源 headers/properties 是有序列表且可能重复 key重复项合并进 map 时最后一个值生效element_at(map, name)查找同样为最后值生效key 不存在时返回NULL。值作为原始字节原样存入 VARCHAR——不做 UTF-8 校验或解码。示例同时加载订单载荷字段order_id、源主题、分区、偏移量、转换为DATETIME的消息时间戳以及trace-idheaderCREATE TABLE example_db.orders_with_meta ( order_id BIGINT, src_topic VARCHAR(256), src_partition INT, src_offset BIGINT, msg_time DATETIME, trace_id VARCHAR(128) ) ENGINE OLAP DUPLICATE KEY(order_id) DISTRIBUTED BY HASH(order_id); CREATE ROUTINE LOAD example_db.orders_with_meta_job ON orders_with_meta INCLUDE METADATA ( TOPIC AS m_topic, PARTITION AS m_partition, OFFSET AS m_offset, TIMESTAMP_MS AS m_timestamp, HEADERS AS m_headers ), COLUMNS ( order_id, src_topic m_topic, src_partition m_partition, src_offset m_offset, msg_time from_unixtime(m_timestamp / 1000), trace_id element_at(m_headers, trace-id) ) PROPERTIES ( format json, jsonpaths [\$.order_id\] ) FROM KAFKA ( kafka_broker_list kafka_broker1_ip:kafka_broker1_port,..., kafka_topic topic_orders, property.kafka_default_offsets OFFSET_BEGINNING );元数据别名可以用于表达式中如上例的from_unixtime(m_timestamp / 1000)jsonpaths中只列出载荷列。查看加载作业与任务查看加载作业执行 SHOW ROUTINE LOAD 语句查看作业example_tbl2_ordertest2的状态。StarRocks 会返回执行状态State、统计信息包括消费的总行数、加载的总行数Statistics以及作业进度Progress。MySQL [example_db] SHOW ROUTINE LOAD FOR example_tbl2_ordertest2 \G *************************** 1. row *************************** Id: 63013 Name: example_tbl2_ordertest2 CreateTime: 2022-08-10 17:09:00 PauseTime: NULL EndTime: NULL DbName: default_cluster:example_db TableName: example_tbl2 State: RUNNING DataSourceType: KAFKA CurrentTaskNum: 3 JobProperties: {partitions:*,partial_update:false,columnToColumnExpr:commodity_id,customer_name,country,pay_time,pay_dtfrom_unixtime(pay_time, %Y%m%d),price,maxBatchIntervalS:20,whereExpr:*,dataFormat:json,timezone:Asia/Shanghai,format:json,json_root:,strict_mode:false,jsonpaths:[\$.commodity_id\,\$.customer_name\,\$.country\,\$.pay_time\,\$.price\],desireTaskConcurrentNum:3,maxErrorNum:0,strip_outer_array:false,currentTaskConcurrentNum:3,maxBatchRows:200000} DataSourceProperties: {topic:ordertest2,currentKafkaPartitions:0,1,2,3,4,brokerList:kafka_broker1_ip:kafka_broker1_port,kafka_broker2_ip:kafka_broker2_port} CustomProperties: {kafka_default_offsets:OFFSET_BEGINNING} Statistic: {receivedBytes:230,errorRows:0,committedTaskNum:1,loadedRows:2,loadRowsRate:0,abortedTaskNum:0,totalRows:2,unselectedRows:0,receivedBytesRate:0,taskExecuteTimeMs:522} Progress: {0:1,1:OFFSET_ZERO,2:OFFSET_ZERO,3:OFFSET_ZERO,4:OFFSET_ZERO} ReasonOfStateChanged: ErrorLogUrls: OtherMsg:状态异常排查如果作业状态自动变为PAUSED可能是因为错误行数超过了阈值。阈值的设置方法见 CREATE ROUTINE LOAD对应的max_error_num属性默认值为 0源码见 RoutineLoadJob.java。可通过ReasonOfStateChanged与ErrorLogUrls定位问题修复问题后执行 RESUME ROUTINE LOAD 恢复PAUSED作业。如果作业状态为CANCELLED通常是因为作业遇到异常如表已被删除。可通过ReasonOfStateChanged与ErrorLogUrls排查但CANCELLED作业无法恢复。注意无法查看已停止或尚未启动的加载作业。查看加载任务执行 SHOW ROUTINE LOAD TASK 语句查看作业example_tbl2_ordertest2的加载任务例如当前正在运行的任务数、被消费的 Kafka 主题分区与消费进度DataSourceProperties以及对应的 Coordinator BE 节点BeId。MySQL [example_db] SHOW ROUTINE LOAD TASK WHERE JobName example_tbl2_ordertest2 \G *************************** 1. row *************************** TaskId: 18c3a823-d73e-4a64-b9cb-b9eced026753 TxnId: -1 TxnStatus: UNKNOWN JobId: 63013 CreateTime: 2022-08-10 17:09:05 LastScheduledTime: 2022-08-10 17:47:27 ExecuteStartTime: NULL Timeout: 60 BeId: -1 DataSourceProperties: {1:0,4:0} Message: there is no new data in kafka, wait for 20 seconds to schedule again *************************** 2. row *************************** TaskId: f76c97ac-26aa-4b41-8194-a8ba2063eb00 TxnId: -1 TxnStatus: UNKNOWN JobId: 63013 CreateTime: 2022-08-10 17:09:05 LastScheduledTime: 2022-08-10 17:47:26 ExecuteStartTime: NULL Timeout: 60 BeId: -1 DataSourceProperties: {2:0} Message: there is no new data in kafka, wait for 20 seconds to schedule again *************************** 3. row *************************** TaskId: 1a327a34-99f4-4f8d-8014-3cd38db99ec6 TxnId: -1 TxnStatus: UNKNOWN JobId: 63013 CreateTime: 2022-08-10 17:09:26 LastScheduledTime: 2022-08-10 17:47:27 ExecuteStartTime: NULL Timeout: 60 BeId: -1 DataSourceProperties: {0:2,3:0} Message: there is no new data in kafka, wait for 20 seconds to schedule again作业生命周期管理暂停作业执行 PAUSE ROUTINE LOAD 语句暂停加载作业作业状态变为PAUSED但并未停止可以执行 RESUME ROUTINE LOAD 恢复也可以用 SHOW ROUTINE LOAD 查看其状态。PAUSE ROUTINE LOAD FOR example_tbl2_ordertest2;恢复作业执行 RESUME ROUTINE LOAD 语句恢复被暂停的加载作业。作业状态会先短暂变为NEED_SCHEDULE因为作业正在被重新调度然后变为RUNNING。这与源码中RoutineLoadScheduler驱动NEED_SCHEDULE - RUNNING状态迁移的设计一致RoutineLoadScheduler.java。RESUME ROUTINE LOAD FOR example_tbl2_ordertest2;修改作业修改作业前必须先使用 PAUSE ROUTINE LOAD 暂停作业然后执行 ALTER ROUTINE LOAD。修改完成后执行 RESUME ROUTINE LOAD 恢复作业并用 SHOW ROUTINE LOAD 查看状态。假设存活 BE 节点数增加到6需要消费的 Kafka 主题分区变为0,1,2,3,4,5,6,7。如果想提高实际加载任务并发度可以执行以下语句将期望任务并发数desired_concurrent_number增加到6大于或等于存活 BE 节点数并指定新的 Kafka 主题分区与初始 offset注意由于实际任务并发数由多个参数的最小值决定必须确保 FE 动态参数max_routine_load_task_concurrent_num的值大于或等于6。ALTER ROUTINE LOAD FOR example_tbl2_ordertest2 PROPERTIES ( desired_concurrent_number 6 ) FROM kafka ( kafka_partitions 0,1,2,3,4,5,6,7, kafka_offsets OFFSET_BEGINNING,OFFSET_BEGINNING,OFFSET_BEGINNING,OFFSET_BEGINNING,OFFSET_END,OFFSET_END,OFFSET_END,OFFSET_END );停止作业执行 STOP ROUTINE LOAD 语句停止加载作业作业状态变为STOPPED。停止后的作业无法恢复也无法用 SHOW ROUTINE LOAD 查看其状态。STOP ROUTINE LOAD FOR example_tbl2_ordertest2;小结Routine Load 是 StarRocks 面向 Kafka 消息流的内置持续加载方案其作业-任务两级模型FE 管理作业与调度、Coordinator/Executor BE 消费并落盘与 exactly-once 语义保障使其成为实时数仓摄入层的可靠选择。实践要点可归纳为并发调优实际并发数 min(分区数, desired_concurrent_number, 存活 BE 数, max_routine_load_task_concurrent_num)扩容时需同步提升 BE 数与 FE 配置上限格式选择CSV 用COLUMNS TERMINATED BYCOLUMNS映射JSON 用format jsonjsonpathsAvrov3.0.1 起用format avroconfluent.schema.registry.url元数据注入JSON/Avro 场景可用INCLUDE METADATA将 topic/partition/offset/timestamp/key/headers 一并写入目标表便于审计与事件时间处理生命周期暂停PAUSED→ 修改ALTER→ 恢复RUNNING→ 停止STOPPED的链路中只有PAUSED状态可恢复STOPPED/CANCELLED均不可恢复。如需深入了解各语句的完整参数请参阅 CREATE ROUTINE LOAD 及其配套的 ALTER、PAUSE、RESUME、STOP、SHOW 与 SHOW ROUTINE LOAD TASK 文档快速上手可直接参考 Quick Start: Routine Load。【免费下载链接】starrocksThe worlds fastest open query engine for sub-second analytics both on and off the data lakehouse. With the flexibility to support nearly any scenario, StarRocks provides best-in-class performance for multi-dimensional analytics, real-time analytics, and ad-hoc queries. A Linux Foundation project.项目地址: https://gitcode.com/GitHub_Trending/st/starrocks创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考