尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

DataX 分区同步实战:基于 datax-web 的动态分区参数配置与调度实现

DataX 分区同步实战:基于 datax-web 的动态分区参数配置与调度实现 数据集成数据同步任务调度后端【免费下载链接】datax-webDataX集成可视化页面选择数据源即可一键生成数据同步任务支持RDBMS、Hive、HBase、ClickHouse、MongoDB等数据源批量创建RDBMS数据同步任务集成开源调度系统支持分布式、增量同步数据、实时查看运行日志、监控执行器资源、KILL运行进程、数据源信息加密等。项目地址https://gitcode.com/gh_mirrors/da/datax-web点击查看免费下载本文围绕 datax-web 仓库中的分区同步方案文档partition-synchronization.md完整讲解 HDFS/Hive 分区表通过 DataX 同步到 ClickHouse 等目标端的实现方式从 DataX Json 中如何用${p_data_day}动态参数代替分区目录、datax.py命令行如何通过-D注入分区值到 datax-web 调度平台如何自动计算当前时间 ± N 天生成分区参数并下发执行。读完本文你将能够独立配置一条按天分区自动同步的 DataX 任务并理解其背后的源码执行链路。一、为什么分区同步需要动态参数DataX 的hdfsreader在读取 Hive/HDFS 分区表时需要显式指定物理路径例如/user/gsbdc/dbdatas/olsd/bns/gsods_rpt_qq/poi/p_data_day2018-05-14/*但hdfsreader本身无法感知分区信息每次同步都要把具体的分区值写死在路径里。而分区表几乎总是按天、按小时滚动产生新分区写死路径意味着每天都要人工改 Json。解决思路即本仓库文档给出的方案是通过 DataX 的-D动态参数把分区值在运行时注入到 Json 中让 Json 模板保持不变、只替换分区值。这一机制在 datax-web 中被抽象为三种增量/动态模式之一对应源码 IncrementTypeEnum.java 中的定义TIME(2, 时间), ID(1, 自增主键), PARTITION(3, HIVE分区);其中PARTITIONHIVE 分区就是本文的核心主题它专门用于按分区值动态同步的场景。二、DataX Json 完整配置样例以下 Json 完整继承自原文档实现 HDFS 分区目录读取 → ClickHouse 写入其中第 6 列通过value: ${p_data_day}将分区值作为一列数据参与同步{ job: { setting: { speed: { channel: 3, byte: 1048576 }, errorLimit: { record: 0, percentage: 0.02 } }, content: [ { reader: { name: hdfsreader, parameter: { hadoopConfig: { dfs.nameservices: nameservice1, dfs.ha.namenodes.nameservice1: cdh201.qq.org,cdh202.qq.org, dfs.namenode.rpc-address.nameservice1.cdh201.qq.org: cdh201.qq.org:8020, dfs.namenode.rpc-address.nameservice1.cdh202.qq.org: cdh202.qq.org:8020, dfs.client.failover.proxy.provider.nameservice1: org.apache.hadoop.hdfs.server.namenode.ha.ConfiguredFailoverProxyProvider }, path: /user/gsbdc/dbdatas/olsd/bns/gsods_rpt_qq/poi/p_data_day2018-05-14/*, haveKerberos: true, kerberosPrincipal: biqq.ORG, defaultFS: hdfs://nameservice1, kerberosKeytabFilePath: /app/soft/datax/job/bi.keytab, fileType: text, fieldDelimiter: \u0001, column: [ { index: 0, type: string }, { index: 1, type: string }, { index: 2, type: string }, { index: 3, type: string }, { index: 4, type: string }, { value: ${p_data_day}, type: string } ] } }, writer: { name: clickhousewriter, parameter: { username: s, password: s, column: [ id, address, p_name, c_name, d_name, p_data_day ], connection: [ { table: [poi], jdbcUrl: jdbc:clickhouse://192.168.1.1:18123/test } ] } } } ] } }配置要点解读配置块说明setting.speed并发通道数channel3与单通道限速byte1048576 字节/秒setting.errorLimit容错上限单条记录错误数 record0、错误率 percentage0.022%超过即终止任务reader.hadoopConfigNameNode HA 相关配置nameservice、namenode 地址、failover proxy provider按你的 Hadoop 集群实际调整reader.haveKerberos / kerberosPrincipal / kerberosKeytabFilePath开启 Kerberos 认证时的认证主体与 keytab 路径无需认证的环境可去掉reader.column前 5 列按index从文本中取值第 6 列不取文件内容而是用value接收动态参数reader.fieldDelimiter字段分隔符Hive 表默认\u0001writer.column与 reader 输出列一一对应最后一个是分区字段p_data_day值得说明的是path中的p_data_day2018-05-14与column中的${p_data_day}是两处独立的存在——前者决定读取哪个分区目录后者决定把分区值作为一列数据写出去。两者都要正确配置缺一不可。三、reader 分区信息的配置方式原文档明确指出DataX hdfsreader 无法获取分区信息必须通过动态参数指定。reader 中分区信息的配置即 Json 的column数组里使用value键来代替index{ value: ${p_data_day}, type: string }这里${...}是 DataX 动态参数的固定格式p_data_day是参数名。DataX 在启动时会把命令行-D传入的键值对替换进 Json 中所有同名的${}占位符。因此要求Json 中占位符名字p_data_day与命令行-D的参数名必须完全一致该列的类型type建议统一为string与分区值文本保持兼容。四、Python 命令行执行方式Json 准备好后通过 DataX 自带的 Python 启动脚本执行python /app/soft/datax/bin/datax.py -p -Dp_data_day2020-06-20 /app/soft/datax/job/hive2clickhouse.json-p是 DataX 的参数注入开关其后的双引号内是-D参数名参数值形式一个任务可注入多个参数例如-p -Dp_data_day2020-06-20 -Dotherxxx多个-D之间以空格分隔注意命令中的p_data_day分区字段要和 reader 中配置的value变量名称一致否则 DataX 无法完成替换会直接按字面量${p_data_day}解析而报错或写出错误数据。五、DataX Web 中的动态传参配置在 datax-web 平台中分区值不必每次手动填写而是由调度平台在任务触发时自动计算。原文档描述的机制为配置定时任务任务执行时获取当前时间及用户选择的当前时间 ± N 天计算得到动态参数的值。也就是说页面上的分区信息配置本质上由三段组成对应源码 BuildCommand.java 中的buildPartition解析逻辑private static String buildPartition(ListString partitionInfo) { String field partitionInfo.get(0); // 分区字段名如 p_data_day int timeOffset Integer.parseInt(partitionInfo.get(1)); // 相对当天的偏移天数可为负数 String timeFormat partitionInfo.get(2); // 时间格式化模板如 yyyy-MM-dd String partitionTime DateUtil.format(DateUtil.addDays(new Date(), timeOffset), timeFormat); return field Constants.EQUAL partitionTime; }三段以英文逗号分隔最终被拼装成-Dpartition分区字段日期追加到命令行参数中见 DataXConstant.java 中的PARAMS_CM_V_PT -Dpartition%s。页面配置示例以每天凌晨同步前一天分区为例在任务管理页面的增量/分区配置区域中配置项示例值说明增量方式HIVE 分区对应IncrementTypeEnum.PARTITIONcode3分区信息p_data_day,-1,yyyy-MM-dd分区字段名 偏移天数 时间格式逗号分隔调度表达式Cron0 0 1 * * ?每天凌晨 1 点触发任务在2020-06-20 01:00:00触发时平台计算2020-06-20 (-1) 天 2020-06-19按yyyy-MM-dd格式化后生成参数-Dpartitionp_data_day2020-06-19。若偏移天数为0则生成当天的分区值。负偏移常用于补昨天数据的 T1 场景正偏移可用于预生成明天分区。时间格式的取值约束时间格式直接复用平台内置的日期格式集合定义在 DateFormatUtils.javapublic static final String DATE_FORMAT yyyy/MM/dd; public static final String DATETIME_FORMAT yyyy/MM/dd HH:mm:ss; public static final String TIME_FORMAT HH:mm:ss; public static final String TIMESTAMP Timestamp;使用时需保证格式化结果与目标分区目录的命名规则一致例如分区目录是p_data_day2020-06-19则时间格式必须填yyyy-MM-dd不能填yyyy/MM/dd。六、源码级执行链路从调度到命令拼装datax-web 的自动计算分区值并非只在页面上生效而是贯穿了调度端与执行端两个模块其完整调用链如下调度端收集分区配置JobTrigger.java 在触发任务时读取任务的incrementType与partitionInfo当IncrementTypeEnum.PARTITION.getCode() incrementType时把partitionInfo即页面填写的p_data_day,-1,yyyy-MM-dd设置进TriggerParam的partitionInfo字段} else if (IncrementTypeEnum.PARTITION.getCode() incrementType) { triggerParam.setPartitionInfo(jobInfo.getPartitionInfo()); }远程下发TriggerParam见 TriggerParam.java作为 RPC 参数携带partitionInfo、replaceParam、jvmParam、startId/endId、startTime/triggerTime等字段通过执行器路由策略下发到对应执行器节点。执行端拼装命令BuildCommand.java 的buildDataXParam负责把上述字段拼成datax.py的完整命令行参数if (incrementType ! null IncrementTypeEnum.PARTITION.getCode() incrementType) { if (StringUtils.isNotBlank(partitionStr)) { ListString partitionInfo Arrays.asList(partitionStr.split(SPLIT_COMMA)); if (doc.length() 0) doc.append(SPLIT_SPACE); doc.append(PARAMS_CM).append(TRANSFORM_QUOTES) .append(String.format(PARAMS_CM_V_PT, buildPartition(partitionInfo))) .append(TRANSFORM_QUOTES); } }最终命令形态拼装结果形如python /app/soft/datax/bin/datax.py -p -Dpartitionp_data_day2020-06-19 /tmp/datax/job/xxx.json其中-j -Xms2G -Xmx2G之类的 JVM 参数JVM_CM -j也会在存在时被先行拼入。整个命令还会写入任务日志JobLogger.log(------------------Command parameters: doc)便于在 datax-web 的运行日志中核对实际生效的分区值。至此可以清楚看到页面填写的分区字段 偏移天数 时间格式在执行时被替换为具体的分区日期再注入 DataX Json 的${partition}或与分区字段同名的占位符。这正是原文档所述获取当前时间及用户选择的当前时间 ± N 天计算得到动态参数值的落地实现。七、相关动态参数模式的对比与衔接分区同步与 datax-web 另外两种动态参数模式共用同一套命令拼装框架理解它们的差异有助于在配置时正确选择增量方式详细配置见 increment-desc.md 与 partition-dynamic-param.md增量方式页面典型配置生成的命令行参数典型应用时间增量TIME辅助参数-DlastTime%s -DcurrentTime%s时间格式选 Timestamp 或yyyy/MM/dd等-p -DlastTime1572537600 -DcurrentTime1579317145按 updateTime 抽取变更数据第一次从页面输入的增量开始时间起跑成功后更新为上次触发时间主键自增ID辅助参数-DstartId%s -DendId%s配置主键字段与增量初始 ID-p -DstartId100 -DendId2000按自增主键分段抽取endId 为本次触发时表内 max(id)HIVE 分区PARTITION分区信息字段,偏移天数,时间格式-p -Dpartitionp_data_day2020-06-19按分区目录做全量对账或分区级同步其中时间增量与分区同步可以组合使用在 partition-dynamic-param.md 的示例中reader 用${lastTime}、${currentTime}圈定增量数据范围writer 的path用${partition}指定写入分区最终拼接结果为-p -DlastTime1572537600 -DcurrentTime1579317145 -Dpartitiondatety2020-01-18——即增量抽数据、动态落分区的完整形态。八、常见问题与注意事项参数名必须一致命令行-D后面的名字、Json 中${}占位符的名字必须完全相同包括大小写。这是分区同步以及时间/ID 增量最常踩的坑。%s占位符格式时间增量的辅助参数-DlastTime%s -DcurrentTime%s中%s是平台用于替换实际时间的占位符必须完整保留且格式完全一致两个-D之间保留且仅保留一个空格。分区偏移的正负-1表示昨天T1 场景常用0表示当天取值要与分区目录实际存在的日期对齐避免读到不存在的目录导致任务失败。时间格式要与分区目录一致分区目录为p_data_day2020-06-19时时间格式必须填yyyy-MM-ddHive 分区若用其他粒度月、小时请选用对应的格式模板可参考 DateFormatUtils.java 内置集合或自行扩展。首次全量与失败语义时间增量/主键增量的起点值在任务成功执行后才会更新为上次触发时间/最大 ID任务失败不会更新可安全重跑。Kerberos 环境hdfsreader 配置了haveKerberostrue时需确保执行器所在节点拥有可读的 keytab 文件如/app/soft/datax/job/bi.keytab否则认证失败。运行日志核对任务触发后datax-web 运行日志中会输出Command parameters行可直接核对实际注入的分区值是否正确日志输出逻辑见 BuildCommand.java 的JobLogger.log调用。通过以上配置与原理梳理即可在 datax-web 中搭建一条稳定运行的按分区自动同步任务Json 模板只写一次分区值由平台按当前时间 ± 偏移自动计算注入实现真正的无人值守分区级数据同步。赞分享数据集成数据同步任务调度后端【免费下载链接】datax-webDataX集成可视化页面选择数据源即可一键生成数据同步任务支持RDBMS、Hive、HBase、ClickHouse、MongoDB等数据源批量创建RDBMS数据同步任务集成开源调度系统支持分布式、增量同步数据、实时查看运行日志、监控执行器资源、KILL运行进程、数据源信息加密等。项目地址https://gitcode.com/gh_mirrors/da/datax-web点击查看免费下载相关推荐DataX-Web动态分区同步Hive分区数据的自动化处理终极指南DataX Web动态分区同步Hive分区数据的自动化处理终极指南 在大数据生态系统中 DataX Web 作为一个强大的数据同步工具提供了动态分区同步功数据集成数据同步任务调度后端如何快速实现HBase数据同步DataX-Web完整实战指南如何快速实现HBase数据同步DataX Web完整实战指南 DataX Web是一个基于DataX开发的分布式数据同步工具的Web管理界面专门用于简化大数数据集成数据同步任务调度后端DataX kuduwriter 插件实战指南Kudu 数据写入与建表分区配置详解DataX kuduwriter 插件实战指南Kudu 数据写入与建表分区配置详解 本篇指南围绕阿里云 DataX 开源版本中的 kuduwriter 写入插数据集成批处理ETL大数据后端上一篇Akagi 麻将AI使用指南三步启动实时牌局分析下一篇思源宋体TTF 7字重下载安装与网页字体配置教程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表