拓冰建站拓冰建站
首页 / 资讯中心 / 正文

使用 Routine Load 从 Apache Pulsar 持续导入数据到 StarRocks

使用 Routine Load 从 Apache Pulsar 持续导入数据到 StarRocks【免费下载链接】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 2.5 起Routine Load 支持从 Apache Pulsar 持续消费消息并写入 StarRocks将 Pulsar 这一存算分离架构的消息流平台作为实时数仓的上游数据源。本文以 CSV 格式数据为例完整讲解 Pulsar 相关核心概念、CREATE ROUTINE LOAD的建作业参数、作业与任务的查看与修改方法并结合仓库源码揭示 Pulsar 消费器在 FE/BE 两侧的实现原理与关键配置。读完本文你将掌握一套可直接落地的 Pulsar → StarRocks 持续导入方案。实验特性说明截至当前仓库版本Routine Load 消费 Pulsar 仍属于实验性Experimental能力从 load_from_pulsar.md 顶部的Experimental /标记可以看出。该功能核心链路已经过测试覆盖参见 PulsarRoutineLoadJobTest.java、PulsarUtilTest.java 与 data_consumer_test.cpp但在生产环境大规模使用前建议先在测试集群充分验证。支持的导入数据格式Routine Load 从 Pulsar 集群消费数据时支持以下两种格式CSV使用列分隔符与行分隔符切分消息流。JSON配合jsonpaths、json_root等参数进行字段映射与解析。注意CSV 格式限制CSV 格式的列分隔符仅支持UTF-8 编码、长度不超过 50 字节的字符串。常用的列分隔符包括逗号,、制表符\t和竖线|行分隔符通常为\n。超出 50 字节的分隔符会校验失败。Pulsar 相关核心概念在使用 Routine Load 消费 Pulsar 之前需要先理解三个与参数配置强相关的 Pulsar 概念。Topic主题Topic 是 Pulsar 中消息从生产者传递到消费者的命名通道分为两类分区主题Partitioned topics由多个 Broker 共同处理可承载更高吞吐。一个分区主题在物理上被实现为 N 个内部主题其中 N 是分区数。非分区主题Non-partitioned topics仅由一个 Broker 服务吞吐上限受限于单个 Broker。从源码看BE 侧的PulsarDataConsumer::get_topic_partition()通过 Pulsar C 客户端提供的getPartitionsForTopic()获取主题的全部分区列表见 data_consumer.cppFE 在作业调度时也会通过 BRPC 代理向 BE 获取分区元数据见下文“调度与消费原理”。Message ID消息 ID消息被 BookKeeper 持久化存储后会获得一个消息 ID用于标识消息在 ledger 中的具体位置且在 Pulsar 集群内唯一。与 Kafka 中是一个long型 offset 不同Pulsar 的消息 ID 由四部分组成ledgerId:entryID:partition-index:batch-index由于消息 ID 无法直接从消息体中取出当前 Routine Load 从 Pulsar 导入数据时不支持指定任意初始位置Message ID仅支持从分区的起始earliest或最新latest位置开始消费。这是与 Kafka Routine Load 最显著的行为差异之一也是pulsar_initial_positions只有POSITION_EARLIEST与POSITION_LATEST两个取值的原因。Subscription订阅订阅是一种命名配置规则决定消息如何投递给消费者。Pulsar 支持一个消费者同时订阅多个主题也支持一个主题上存在多个订阅。订阅类型在消费者连接时确定共有四种订阅类型说明exclusive默认只允许一个消费者附着到该订阅仅一个消费者消费消息。shared多个消费者可附着到同一订阅消息按轮询方式分发每条消息只投递给一个消费者。failover多个消费者可附着到同一订阅对非分区主题或分区主题的每个分区会选举一个主消费者接收消息主消费者断开后未确认及后续消息投递给下一个消费者。key_shared多个消费者可附着到同一订阅相同 key 或 ordering key 的消息只投递给同一个消费者。注意目前 Routine Load 使用exclusive订阅类型这一点可结合 BE 侧实现印证——data_consumer.cpp 中的_p_client-subscribe(partition, _subscription, _p_consumer)直接以订阅名创建消费者未指定ConsumerConfiguration中的订阅类型即使用 Pulsar 默认的Exclusive模式。因此同一订阅名在同一时刻只能被一个 Routine Load 作业消费避免多个消费者共享订阅导致消息竞争。创建 Routine Load 作业下面以一个完整的示例说明消费 Pulsar 中 CSV 格式的消息通过创建 Routine Load 作业routine_wiki_edit_1将数据持续写入表routine_wiki_edit。CREATE ROUTINE LOAD load_test.routine_wiki_edit_1 ON routine_wiki_edit COLUMNS TERMINATED BY ,, ROWS TERMINATED BY \n, COLUMNS (order_id, pay_dt, customer_name, nationality, temp_gender, price) WHERE event_time 2022-01-01 00:00:00, PROPERTIES ( desired_concurrent_number 1, max_batch_interval 15000, max_error_number 1000 ) FROM PULSAR ( pulsar_service_url pulsar://localhost:6650, pulsar_topic persistent://tenant/namespace/topic-name, pulsar_subscription load-test, pulsar_partitions load-partition-0,load-partition-1, pulsar_initial_positions POSITION_EARLIEST,POSITION_LATEST, property.auth.token eyJ0eXAiOiJKV1QiLCJhbGciOiJIUJzdWIiOiJqaXV0aWFuY2hlbiJ9.lulGngOC72vE70OW54zcbyw7XdKSOxET94WT_hIqD5Y );创建 Pulsar Routine Load 作业时除data_source_properties之外的绝大多数参数如COLUMNS TERMINATED BY、ROWS TERMINATED BY、COLUMNS、WHERE、desired_concurrent_number、max_batch_interval、max_error_number等与消费 Kafka 时完全一致完整参数说明请参考 CREATE_ROUTINE_LOAD.md。data_source_properties 参数详解FROM PULSAR (...)内与数据源相关的参数如下参数必填说明pulsar_service_url是连接 Pulsar 集群的 URL格式为pulsar://ip:port或pulsar://service:port例如pulsar_service_url pulsar://localhost:6650。pulsar_topic是订阅的主题例如pulsar_topic persistent://tenant/namespace/topic-name。pulsar_subscription是为主题配置的订阅名例如pulsar_subscription my_subscription。pulsar_partitions、pulsar_initial_positions否pulsar_partitions订阅的主题分区列表pulsar_initial_positionspulsar_partitions中各分区对应的初始消费位置数量必须与分区一一对应。有效取值POSITION_EARLIEST默认从分区最早可用消息开始消费POSITION_LATEST从分区最新可用消息开始消费。注意若未指定pulsar_partitions则订阅主题的全部分区若同时指定了pulsar_partitions与property.pulsar_default_initial_position以pulsar_partitions的取值为准若两者都未指定默认从各分区最新可用消息开始消费。示例pulsar_partitions my-partition-0,my-partition-1,my-partition-2,my-partition-3, pulsar_initial_positions POSITION_EARLIEST,POSITION_EARLIEST,POSITION_LATEST,POSITION_LATESTRoutine Load 还支持以下 Pulsar 自定义参数参数必填说明property.pulsar_default_initial_position否未指定pulsar_initial_positions时主题各分区订阅的默认初始位置有效取值与pulsar_initial_positions相同。示例property.pulsar_default_initial_position POSITION_EARLIESTproperty.auth.token否若 Pulsar 启用了基于安全令牌token的客户端认证需提供 token 字符串以验证身份。示例property.auth.token eyJ0eXAiOiJKV1QiLCJhbGciOiJIUJzdWIiOiJqaXV0aWFuY2hlbiJ9.lulGngOC72vE70OW54zcbyw7XdKSOxET94WT_hIqD从源码看参数校验逻辑FE 侧在 CreateRoutineLoadStmt.java 中定义了上述参数的常量名PULSAR_SERVICE_URL_PROPERTY、PULSAR_TOPIC_PROPERTY、PULSAR_SUBSCRIPTION_PROPERTY、PULSAR_PARTITIONS_PROPERTY、PULSAR_INITIAL_POSITIONS_PROPERTY、PULSAR_DEFAULT_INITIAL_POSITION并在分析阶段做了如下校验了解这些校验有助于避免建作业失败pulsar_service_url、pulsar_topic、pulsar_subscription为必填缺失会抛出xxx is a required propertypulsar_service_url需匹配ENDPOINT_REGEX端点正则支持逗号分隔的多个地址pulsar_partitions不能为空字符串且每个分区在创建作业时会通过 BRPC 向 BE 查询主题真实分区列表做分区合法性校验见 PulsarRoutineLoadJob.java 的checkCustomPartition()pulsar_initial_positions的分区数量必须与pulsar_partitions数量一致否则报错Partitions number should be equals to positions number初始位置只接受POSITION_EARLIEST或POSITION_LATEST不区分大小写内部被映射为数值POSITION_LATEST 0、POSITION_EARLIEST 1见 PulsarRoutineLoadJob.java 与 CreateRoutineLoadStmt.java。从源码看消费初始位置的落地BE 侧的 data_consumer.cpp 中PulsarDataConsumer::assign_partition()完成了订阅与 seek调用_p_client-subscribe(partition, _subscription, _p_consumer)以独占订阅方式创建消费者若指定了初始位置则调用_p_consumer.seek(...)POSITION_LATEST对应pulsar::MessageId::latest()POSITION_EARLIEST对应pulsar::MessageId::earliest()两个步骤任一失败都会返回Status::InternalError并携带PAUSE:前缀——这意味着作业会进入 PAUSED 状态可通过SHOW ROUTINE LOAD查看ReasonOfStateChanged字段定位原因。另外注意pulsar_initial_positions只在作业首次调度时生效。PulsarProgress中的partitionToInitialPosition注释明确写着 Initial positions will only be used at first schedule且每次任务提交后该分区对应的初始位置记录会被移除见 PulsarProgress.java。也就是说作业正常运行后的断点续传依赖于订阅内已确认acknowledged的消息位置而非重复执行 seek。查看导入作业与导入任务查看导入作业SHOW ROUTINE LOAD执行 SHOW_ROUTINE_LOAD.md 中的语句可查看作业routine_wiki_edit_1的状态。检查消费 Pulsar 的 Routine Load 作业时除progress外的大部分返回参数与消费 Kafka 时相同progress在这里表示backlog即分区中未确认unacked的消息数量。MySQL [load_test] SHOW ROUTINE LOAD for routine_wiki_edit_1 \G *************************** 1. row *************************** Id: 10142 Name: routine_wiki_edit_1 CreateTime: 2022-06-29 14:52:55 PauseTime: 2022-06-29 17:33:53 EndTime: NULL DbName: default_cluster:test_pulsar TableName: test1 State: PAUSED DataSourceType: PULSAR CurrentTaskNum: 0 JobProperties: {partitions:*,rowDelimiter:\n,partial_update:false,columnToColumnExpr:*,maxBatchIntervalS:10,whereExpr:*,timezone:Asia/Shanghai,format:csv,columnSeparator:,,json_root:,strict_mode:false,jsonpaths:,desireTaskConcurrentNum:3,maxErrorNum:10,strip_outer_array:false,currentTaskConcurrentNum:0,maxBatchRows:200000} DataSourceProperties: {serviceUrl:pulsar://localhost:6650,currentPulsarPartitions:my-partition-0,my-partition-1,topic:persistent://tenant/namespace/topic-name,subscription:load-test} CustomProperties: {auth.token:eyJ0eXAiOiJKV1QiLCJhbGciOiJIUJzdWIiOiJqaXV0aWFuY2hlbiJ9.lulGngOC72vE70OW54zcbyw7XdKSOxET94WT_hIqD} Statistic: {receivedBytes:5480943882,errorRows:0,committedTaskNum:696,loadedRows:66243440,loadRowsRate:29000,abortedTaskNum:0,totalRows:66243440,unselectedRows:0,receivedBytesRate:2400000,taskExecuteTimeMs:2283166} Progress: {my-partition-0(backlog): 100,my-partition-1(backlog): 0} ReasonOfStateChanged: ErrorLogUrls: OtherMsg: 1 row in set (0.00 sec)对输出中几个 Pulsar 相关字段的解读DataSourceType: PULSAR数据源类型为 PulsarDataSourceProperties包含serviceUrl服务地址、topic主题、subscription订阅名与currentPulsarPartitions当前实际消费的分区。该 JSON 由 FE 的dataSourcePropertiesJsonToString()方法生成currentPulsarPartitions会被排序后输出见 PulsarRoutineLoadJob.javaCustomProperties自定义属性如auth.token会原样回显Statistic统计信息包括接收字节数receivedBytes、导入行数loadedRows、总行数totalRows、错误行数errorRows、提交任务数committedTaskNum、中止任务数abortedTaskNum、吞吐loadRowsRate/receivedBytesRate等其 JSON 结构定义在getStatistic()方法中见 PulsarRoutineLoadJob.javaProgress形如{my-partition-0(backlog): 100,my-partition-1(backlog): 0}即各分区的未确认消息积压量。BE 侧通过_p_consumer.getBrokerConsumerStats()获取MsgBacklog指标见 data_consumer.cppFE 侧通过 PulsarUtil.java 的getBacklogNums()批量查询后汇总展示。查看导入任务SHOW ROUTINE LOAD TASK执行 SHOW_ROUTINE_LOAD_TASK.md 中的语句可查看作业routine_wiki_edit_1的各个导入任务包括当前运行任务数、消费的分区、消费进度DataSourceProperties以及对应的 Coordinator BE 节点BeId。MySQL [example_db] SHOW ROUTINE LOAD TASK WHERE JobName routine_wiki_edit_1 \G任务层面的调度逻辑可以参考 FE 的divideRoutineLoadJob()作业被调度时会将currentPulsarPartitions按j % currentConcurrentTaskNum i的方式轮询均分到各并发任务中并为每个任务记录分区及首次调度的初始位置见 PulsarRoutineLoadJob.java。并发任务数currentTaskConcurrentNum取分区数、期望并发数desireTaskConcurrentNum与存活 BE/CN 数量的最小值并受max_routine_load_task_concurrent_num配置约束见 PulsarRoutineLoadJob.java。修改导入作业修改导入作业前必须先使用 PAUSE_ROUTINE_LOAD.md 暂停作业再执行 ALTER_ROUTINE_LOAD.md 进行修改修改完成后用 RESUME_ROUTINE_LOAD.md 恢复作业并通过 SHOW_ROUTINE_LOAD.md 查看状态。Routine Load 消费 Pulsar 时除data_source_properties外的大部分参数修改方式与消费 Kafka 时相同。需要特别留意以下几点限制在data_source_properties相关参数中目前仅支持修改pulsar_partitions、pulsar_initial_positions以及自定义参数property.pulsar_default_initial_position和property.auth.tokenpulsar_service_url、pulsar_topic、pulsar_subscription不可修改若要修改消费的分区及其对应的初始位置创建作业时必须通过pulsar_partitions显式指定分区且只能修改这些已指定分区的初始位置pulsar_initial_positions若创建作业时只指定了主题pulsar_topic而未指定分区pulsar_partitions则可通过pulsar_default_initial_position修改主题下所有分区的起始位置。从源码看这些限制在 FE 的modifyDataSourceProperties()中强制执行修改分区初始位置时会逐一校验分区是否在创建语句中指定否则抛出The partition xxx is not specified in the create statement见 PulsarRoutineLoadJob.java。另外当作业被暂停且满足自动调度条件时调度器会将作业从PAUSED恢复到NEED_SCHEDULE状态见 PulsarRoutineLoadJob.java 的applyPausedAutoSchedule()配合unprotectUpdateProgress()将新分区补充默认初始位置见 PulsarProgress.java。消费链路与关键配置BE 侧实现原理了解 BE 侧的消费实现有助于理解作业行为与调优方向。Pulsar 消费器实现在 data_consumer.h 的PulsarDataConsumer类中BE 编译时通过#ifndef __APPLE__控制macOS 下不参与编译其消费主循环位于group_consume()见 data_consumer.cpp以1 秒超时receive(..., 1000 /* timeout, ms */)循环拉取单条消息消息成功入队后放入TimedBlockingQueuepulsar::Message*交给下游 Stream Load 管道解析写入ResultTimeout视为正常空转上游暂时无数据其他错误则中止任务消息消费完成后通过acknowledge_cumulative()调用_p_consumer.acknowledgeCumulative(messageId)做累积确认见 data_consumer.cpp这也是progress中 backlog 不断被消耗的机制来源。init()阶段见 data_consumer.cpp还会做两件重要的事遍历自定义属性仅识别auth.token并通过pulsar::AuthToken::createWithToken()配置客户端认证其他属性打印Config xxx not supported for now警告后忽略使用FileLoggerFactory将 Pulsar C 客户端日志输出到${sys_log_dir}/pulsar-cpp-client.log日志级别由配置项控制。相关的 BE 配置项定义在 config.h 中可在be.conf中按需调整配置项默认值说明max_pulsar_consumer_num_per_group10单个数据消费者组中允许的最大 Pulsar 消费者数量用于 Routine Load见 config.h。routine_load_pulsar_timeout_second10Pulsar 请求超时时间秒FE 通过 BRPC 向 BE 代理查询分区与 backlog 时使用见 config.h消费方见 PulsarUtil.java。pulsar_client_log_level2Pulsar C 客户端日志级别0DEBUG1INFO2WARN3ERROR默认 WARN见 config.h。停止导入作业当不再需要从 Pulsar 持续导入时可执行 STOP_ROUTINE_LOAD.md 中的语句停止作业。停止后作业无法恢复如需继续导入需重新创建作业。最佳实践与注意事项结合文档与源码给出以下几点落地建议订阅名全局唯一性由于 Routine Load 使用 exclusive 订阅同一订阅名只能被一个作业消费。若要重启导入流水线务必先STOP旧作业或等待其释放订阅避免subscribe阶段冲突。初始位置只生效一次pulsar_initial_positions与property.pulsar_default_initial_position只在首次调度时通过seek定位。作业运行后若需要回溯消费应使用PAUSE→ALTER ROUTINE LOAD修改初始位置 →RESUME的流程。分区增删的自动感知作业运行时FE 调度器会周期性刷新分区列表refreshPartitionsIfNeeded()见 PulsarRoutineLoadJob.java主题新增分区会被纳入消费、移除分区会被剔除并触发重新调度若创建作业时通过pulsar_partitions固定了分区列表则不会触发 Broker RPC仅使用固定列表。合理设置并发与批处理desired_concurrent_number决定任务并发度实际还受分区数与存活 BE/CN 数约束max_batch_interval控制攒批时长、max_error_number控制容错上限。消息量大时可适当调高并发与批量参数以提升吞吐。CSV 分隔符注意列分隔符必须是 UTF-8 编码、不超过 50 字节的字符串且要确保与消息内容不冲突消息字段含分隔符时应优先选择 JSON 格式或使用不易冲突的分隔符。监控 backlog 指标SHOW ROUTINE LOAD的Progress字段直接反映各分区未确认消息积压。若 backlog 持续增长说明消费速度低于生产速度可结合Statistic中的loadRowsRate与receivedBytesRate判断瓶颈再决定是否扩容 BE/CN 或调高并发。参考文档CREATE ROUTINE LOADSHOW ROUTINE LOADSHOW ROUTINE LOAD TASKPAUSE ROUTINE LOADRESUME ROUTINE LOADALTER ROUTINE LOADSTOP 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),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门