十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

SeaTunnel DataHub Sink 连接器:将数据写入阿里云 DataHub 的配置与源码解析

SeaTunnel DataHub Sink 连接器:将数据写入阿里云 DataHub 的配置与源码解析 SeaTunnel DataHub Sink 连接器将数据写入阿里云 DataHub 的配置与源码解析【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文基于 SeaTunnel 官方文档 DataHub Sink 编写系统讲解 SeaTunnel 如何通过 DataHub sink 连接器把数据写入阿里云 DataHub包括完整的 Sink 配置参数、单表写入与多表路由${table}占位符的实战配置、字段按名映射的机制约束并结合 connector-datahub 模块源码剖析客户端构建、LZ4 压缩写入与有界重试的实现细节。读完后你可以直接编写可在 Spark、Flink 或 Zeta 引擎运行的 DataHub 写入作业并理解每个参数在底层代码中如何生效。支持的引擎与连接器定位DataHub sink 连接器支持以下执行引擎SparkFlinkSeaTunnel Zeta其核心职责是把 SeaTunnel 行SeaTunnelRow写入阿里云 DataHub。连接器同时支持单表写入和多表写入在多表作业中可以通过topic中的${table}占位符把来自不同输入表的数据路由到不同的 DataHub topic。关键特性特性是否支持说明exactly-once✗不做精确一次写入见下方 FAQ 中的语义说明CDC✗不支持 CDC 变更流写入多表写入✓通过${table}占位符按表名路由到不同 topictimer flush✗无定时刷盘机制前置条件先创建 Project 和 Topic在运行 SeaTunnel 作业之前必须先在 DataHub 侧完成两件事创建 DataHubproject和topic——连接器不会自动创建目标 topic确保 topic 的 schema 中包含与上游 SeaTunnel schema同名字段——因为 sink 是按字段名写值的writes values by field name字段名对不上就无法正确映射。这一点在后文的源码分析中会得到印证Writer 在写入时通过dataHubClient.getTopic(project, topic).getRecordSchema()拉取远端 topic 的RecordSchema再用上游行类型的字段名逐个setField因此字段名匹配是硬性约束。Sink 配置参数总览参数类型必填默认值说明endpointstring是-DataHub 服务端点通常以http/https开头accessIdstring是-访问 DataHub 使用的阿里云 AccessKey IDaccessKeystring是-访问 DataHub 使用的阿里云 AccessKey Secretprojectstring是-DataHub project 名称topicstring是-DataHub topic 名称多表作业中支持${table}占位符timeoutint否3000客户端最大连接超时时间毫秒retryTimesint否3写入记录失败时的最大重试次数common-optionsconfig否-Sink 插件公共参数见 Sink Common Options这些参数的定义与文档表格完全一致可对照源码 DataHubSinkOptionsENDPOINT、ACCESS_ID、ACCESS_KEY、PROJECT、TOPIC均为noDefaultValue()必填且无默认值TIMEOUT默认3000RETRY_TIMES默认3。topic 与占位符单表作业topic写死为目标 topic 名称必填缺失或为空会在作业启动时校验失败后文有测试用例佐证多表作业topic可含占位符${table}每个输入表会被路由到与表名同名的 topic${table_name}仅作为废弃的兼容别名保留新作业应统一使用${table}。任务示例示例一单表写入单 Topic一个简单的批处理作业把 FakeSource 的数据写入单个 DataHub topicenv { parallelism 1 job.mode BATCH } source { FakeSource { plugin_output fake schema { fields { name string age int } } } } sink { DataHub { endpoint https://datahub.example.aliyuncs.com accessId your-access-id accessKey your-access-key project demo_project topic user_topic timeout 3000 retryTimes 3 } }示例二多张输入表路由到同名 Topic当上游 source 输出多张表时把topic配置为含${table}占位符的值每张输入表即路由到与其同名的 topic。以下示例中 FakeSource 输出users100 行与orders200 行两张表分别写入 DataHub 的users和orderstopic——两个 topic 需提前在 DataHub 中创建且字段名name/age、order_id/amount与上游 schema 一致env { parallelism 1 job.mode BATCH } source { FakeSource { plugin_output fake tables_configs [ { row.num 100 schema { table users fields { name string age int } } }, { row.num 200 schema { table orders fields { order_id int amount decimal(10, 2) } } } ] } } sink { DataHub { endpoint https://datahub.example.aliyuncs.com accessId your-access-id accessKey your-access-key project demo_project topic ${table} timeout 3000 retryTimes 3 } }公共参数DataHub sink 遵循 Sink Common Options例如在多路 pipeline 中通过plugin_input指定数据来源此时上游必须设置plugin_output。多表写入场景下工厂选项规则中还显式声明了multi_table_sink_replica作为可选公共参数——这在 DataHubSinkFactory 的optionRule()中可以直接看到.optional(SinkConnectorCommonOptions.MULTI_TABLE_SINK_REPLICA)一行。源码解析参数校验、客户端构建与写入链路工厂与必填项校验DataHubSinkFactory 实现了TableSinkFactory接口插件标识符为DataHub。其optionRule()声明了 5 个必填项并全部附加了Conditions.notBlank条件——这意味着必填项不仅不能缺失空字符串或纯空白字符同样会被拒绝return OptionRule.builder() .required(ENDPOINT, Conditions.notBlank(ENDPOINT)) .required(ACCESS_ID, Conditions.notBlank(ACCESS_ID)) .required(ACCESS_KEY, Conditions.notBlank(ACCESS_KEY)) .required(PROJECT, Conditions.notBlank(PROJECT)) .required(TOPIC, Conditions.notBlank(TOPIC)) .optional(SinkConnectorCommonOptions.MULTI_TABLE_SINK_REPLICA) .optional(TIMEOUT, RETRY_TIMES) .build();这与单测 DataHubFactoryTest 的断言一一对应testMissingRequiredOptionsRejected逐个移除 5 个必填 key 均期望抛出OptionValidationExceptiontestEmptyRequiredOptionsRejected与testWhitespaceOnlyRequiredOptionsRejected则验证空串和 \t这类纯空白值都会被拒绝。这也解释了 FAQ 中“为什么单表作业topic也是必填”——它在启动期就被规则校验拦下而不是运行时才猜测归属 project。DataHubSink继承AbstractSimpleSink并实现SupportMultiTableSink标记接口见 seatunnel-api 中的空标记定义引擎据此识别其多表写入能力createWriter会把 8 个构造参数原样透传给DataHubWriter。客户端构建timeout 与 LZ4 压缩DataHubWriter 的构造函数揭示了各参数的底层落点this.dataHubClient DatahubClientBuilder.newBuilder() .setDatahubConfig( new DatahubConfig( endpoint, new AliyunAccount(accessId, accessKey), true)) .setHttpConfig( new HttpConfig() .setCompressType(HttpConfig.CompressType.LZ4) .setConnTimeout(timeout)) .build();可以看到endpoint/accessId/accessKey通过DatahubConfigAliyunAccount注入官方 SDKcom.aliyun.datahub:aliyun-sdk-datahubpom.xml 中版本为2.19.0-publictimeout对应HttpConfig.setConnTimeout即HTTP 连接超时而非读写总超时压缩类型被固定为 LZ4无需也不可通过作业配置调整。写入路径按字段名映射 有界重试write(SeaTunnelRow)是理解“字段按名写入”的关键String[] fieldNames seaTunnelRowType.getFieldNames(); Object[] fields element.getFields(); ListRecordEntry recordEntries new ArrayList(); RecordSchema recordSchema dataHubClient.getTopic(project, topic).getRecordSchema(); for (int i 0; i fieldNames.length; i) { TupleRecordData data new TupleRecordData(recordSchema); data.setField(fieldNames[i], fields[i]); // 按上游字段名写值 ... } PutRecordsResult result dataHubClient.putRecords(project, topic, recordEntries);每行上游数据被拆解为逐字段的TupleRecordData以远端 topic 的RecordSchema为基准、用上游字段名setField赋値再一次putRecords批量提交——这正是文档强调“topic schema 字段名必须与上游 schema 一致”的代码依据。失败处理分两层部分记录失败putRecords返回后检查getFailedRecordCount()大于 0 则进入私有方法retry(records, retryTimes, ...)按retryTimes默认 3做有界重试。从源码结构看该重试方法在首轮重投成功后即break返回后续轮次由retryNums递减控制属于“尽力而为的有界重投”客户端异常DatahubClientException被捕获后记录requestId与message日志。结合这两点可以推断连接器整体是 best-effort 语义不承诺 exactly-once——重试预算耗尽后的失败交由作业层处理若上游可接受 at-least-once 重放应在作业层开启 checkpoint 重放。FAQ继承自官方文档DataHub sink 是否支持 exactly-once 投递不支持。连接器执行的是带界重试的尽力写入retryTimes默认3。超出重试预算的失败会上抛给作业而不是被静默吞掉如果上游允许 at-least-once 重放请在作业层启用 checkpoint。多表路由是如何工作的当上游 source 输出多于一张表时把topic设为包含${table}占位符的值例如topic ${table}每张输入表即被路由到与其同名的 DataHub topic。${table_name}仍被识别为废弃别名新作业应优先使用${table}。连接器不会自动创建目标 topic——请提前在 DataHub project 中建好 topic并确保其 schema 字段名与上游 SeaTunnel schema 一致。为什么单表作业中topic也是必填字段DataHub 写入是 schema 绑定的sink 按照目标 topic 声明的字段名序列化每一行。缺失或为空的topic会让连接器无处可写因此它在作业启动时即被校验拦截见工厂optionRule与对应单测而不是尝试从 project 推断。变更记录Changelog 摘录完整变更历史见 connector-datahub changelog变更版本[Feature][Connector-V2] Make some sink parameters optional for DataHub2.3.11[Feature][Connector-V2] Datahub support multi-table sink2.3.11[improve] datahub sink options2.3.10[Improve][Connector-V2][DataHub] Unified exception for DataHub sink connector change package name of DataHub2.3.0[Improve][Connector-V2][DataHub] Add DataHub Sink Factory2.3.0[Feature][Connector-V2]Support datahub sink2.2.0-beta其中多表写入${table}路由与部分参数可选化均于 2.3.11 版本引入使用旧版本发行包时不具备这些能力请以你所使用的发行版本为准。延伸阅读Sink Common Optionsplugin_input、multi_table_sink_replica等公共参数连接器源码DataHubSink、DataHubSinkFactory、DataHubWriter参数定义DataHubSinkOptions配置校验测试DataHubFactoryTest【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表