
数据集成ETL大数据批处理流处理变更数据捕获【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/GitHub_Trending/se/seatunnel点击查看免费下载导读本文围绕 SeaTunnel 的JsonPath 转换插件Transform展开系统讲解如何利用 JSONPath 表达式从 STRING、BYTES、ARRAY、MAP、ROW 等类型的源字段中提取任意层级的嵌套字段并完成类型转换与异常数据处置。读完本文你将掌握 JsonPath 插件的全部配置项columns、src_field、path、dest_field、dest_type、row_error_handle_way、column_error_handle_way、单字段与批量字段两种提取写法以及 FAIL / SKIP / SKIP_ROW 三种异常处理策略在真实数据管道中的组合用法。插件概述JsonPath 是 SeaTunnel 内置的 V2 转换插件其核心能力是使用 JSONPath 选择器从 JSON 数据中提取字段。在数据集成场景中源端如 Kafka、HTTP、文件传入的常常是一整段嵌套 JSON而下游如数据库、数据仓库需要的是拍平后的明细字段。JsonPath 插件可以把这种「解析 提取 类型转换」的工作在管道内一站式完成无需编写 UDF 或依赖外部解析服务。从源码结构看该插件位于 seatunnel-transforms-v2/src/main/java/org/apache/seatunnel/transform/jsonpath/包含四个核心类类职责JsonPathTransformFactory.java插件工厂注册JsonPath标识、声明选项规则与校验器JsonPathTransformConfig.java配置解析与校验将 HOCON 配置转换为ColumnConfig列表JsonPathTransform.java核心转换逻辑执行 JSONPath 提取与类型转换JsonPathMultiCatalogTransform.java多表MultiCatalog场景的适配入口属性配置JsonPath 插件的顶层属性如下名称类型是否必须默认值columnsArrayYes无row_error_handle_wayEnumNoFAIL其中columns为必填项用于声明每条提取规则row_error_handle_way为整行级别的异常处理策略默认FAIL。通用选项plugin_input、plugin_output、multi_tables、table_match_regex、rule_match_mode等转换插件的公共参数请参考 Transform Plugin 通用选项。从 JsonPathTransformFactory.java 的optionRule()可以看到除必填的columns外插件还注册了multi_tables、table_match_regex、rule_match_mode与row_error_handle_way作为可选选项工厂通过AutoService(Factory.class)自动注册到插件体系。row_error_handle_way [Enum]该选项用于指定当该行发生错误时的处理方式默认值为FAIL。在 TransformCommonOptions.java 中定义FAIL选择FAIL时数据格式错误会阻塞并抛出异常任务终止。SKIP选择SKIP时数据格式错误会跳过该行数据。ROUTE_TO_TABLE源码中该选项还支持ROUTE_TO_TABLE将异常数据路由到指定错误表可通过row_error_handle_way.error_table指定目标表名见 TransformCommonOptions.java。columns [array]columns是提取规则的数组数组内每个元素为一个对象包含以下属性名称类型是否必须默认值src_fieldStringYes无dest_fieldString or ArrayYes无pathString or ArrayYes无dest_typeString or ArrayNoStringcolumn_error_handle_wayEnumNo无src_field要解析的 JSON 源字段。src_field必须是输入表中真实存在的字段名支持以下 SeaTunnel 数据类型STRING直接将字段值作为 JSON 字符串解析BYTES将字节数组按字符串解码后解析ARRAY通过 JSON 序列化后解析MAP通过 JSON 序列化后解析ROW将SeatunnelRow的字段数组序列化为 JSON 后解析这一行为在 JsonPathTransform.java 的doTransform方法中有明确实现对 STRING 直接value.toString()BYTES 走new String((byte[]) value)ARRAY/MAP 走JsonUtils.toJsonString(value)ROW 走JsonUtils.toJsonString(row.getFields())。若源字段类型不在上述范围内会抛出unsupportedDataType异常。此外配置解析阶段JsonPathTransformConfig.java会校验src_field必须存在于输入表 Schema 中否则抛出cannotFindInputFieldError。dest_field使用 JSONPath 后的输出字段。可以是单个字段名当需要从同一个源字段提取多个值时也可以配置为字段名数组与path数组一一对应。dest_type目标字段的类型。可以是单个类型也可以在批量提取时配置为类型数组。单字段提取时如果省略默认使用string该默认值定义在 JsonPathTransformConfig.java 的DEST_TYPE选项中。pathJSONPath。可以是单个 JSONPath 表达式也可以是 JSONPath 表达式数组。底层使用com.jayway.jsonpath.JsonPath作为解析引擎并通过JSON_PATH_CACHEConcurrentHashMap缓存编译后的JsonPath对象避免每条数据重复编译表达式见 JsonPathTransform.java。column_error_handle_way [Enum]该选项用于指定当某一列发生错误时的处理方式优先级高于row_error_handle_wayFAIL选择FAIL时数据格式错误会阻塞并抛出异常。SKIP选择SKIP时数据格式错误会跳过此列数据填充空值。SKIP_ROW选择SKIP_ROW时数据格式错误会跳过此行数据。枚举定义位于 ErrorHandleWay.java包含FAIL、SKIP、SKIP_ROW、ROUTE_TO_TABLE四个取值。读取 JSON 示例假设从源读取到的数据是下面这样的 JSON{ data: { c_string: this is a string, c_boolean: true, c_integer: 42, c_float: 3.14, c_double: 3.14, c_decimal: 10.55, c_date: 2023-10-29, c_datetime: 16:12:43.459, c_array:[item1, item2, item3], c_map_array: [{c_string_1:c_string_1,c_string_2:c_string_2,c_string_3:c_string_3},{c_string_1:c_string_1,c_string_2:c_string_2,c_string_3:c_string_3}] } }假设我们想要使用 JSONPath 提取属性可以这样配置transform { JsonPath { plugin_input fake plugin_output fake1 columns [ { src_field data path $.data.c_string dest_field c1_string }, { src_field data path $.data.c_boolean dest_field c1_boolean dest_type boolean }, { src_field data path $.data.c_integer dest_field c1_integer dest_type int }, { src_field data path $.data.c_float dest_field c1_float dest_type float }, { src_field data path $.data.c_double dest_field c1_double dest_type double }, { src_field data path $.data.c_decimal dest_field c1_decimal dest_type decimal(4,2) }, { src_field data path $.data.c_date dest_field c1_date dest_type date }, { src_field data path $.data.c_datetime dest_field c1_datetime dest_type time }, { src_field data path $.data.c_array dest_field c1_array dest_type arraystring }, { src_field data path $.data.c_map_array dest_field c1_map_array dest_type arraymapstring, string } ] } }批量字段提取使用批量字段提取功能可以用更简洁的数组格式配置实现相同的结果transform { JsonPath { plugin_input fake plugin_output fake1 columns [ { src_field data path [$.data.c_string, $.data.c_boolean, $.data.c_integer, $.data.c_float, $.data.c_double, $.data.c_decimal, $.data.c_date, $.data.c_datetime, $.data.c_array, $.data.c_map_array] dest_field [c1_string, c1_boolean, c1_integer, c1_float, c1_double, c1_decimal, c1_date, c1_datetime, c1_array, c1_map_array] dest_type [string, boolean, int, float, double, decimal(4,2), date, time, arraystring, arraymapstring, string] } ] } }重要提示使用批量字段提取时path、dest_field和dest_type的数组长度必须一致。如果省略dest_typeTransform 只会使用单个默认类型string因此多个输出字段应显式配置dest_type数组。这一约束在配置解析层有强制校验JsonPathTransformConfig.of()中pathArray.length ! destFieldArray.length || pathArray.length ! typeArray.length时会直接抛出TransformExceptionPath, dest_field, and dest_type arrays must have the same length见 JsonPathTransformConfig.java。parseFields方法L181-L204会将单个字符串统一转换为单元素数组因此单字段与批量两种写法共用同一套校验逻辑。那么数据结果表fake1将会像这样datac1_stringc1_booleanc1_integerc1_floatc1_doublec1_decimalc1_datec1_datetimec1_arraytoo much content not to showthis is a stringtrue423.143.1410.552023-10-2916:12:43.459[item1, item2, item3]类型转换细节dest_type通过JsonToRowConverters创建对应类型的转换器见 JsonPathTransform.java支持boolean、int、float、double、decimal(p,s)、date、time、timestamp、array...、map...、row...等 SeaTunnel 类型。测试用例 JsonPathTransformTest.java 验证了多种日期格式1990/05/20、2024-01-15、2024/01/15 10:30:00可被正确解析为LocalDate/LocalDateTime闰年日期2024-02-29也可正常处理。读取 SeatunnelRow 示例假设数据行中的一列的类型是SeatunnelRow列的名称为colSeatunnelRow(col)othernameage....a18....JsonPath 转换会将 SeatunnelRow 的值转换为一个 JSON 数组然后用下标路径$[0]、$[1]提取子字段transform { JsonPath { plugin_input fake plugin_output fake1 row_error_handle_way FAIL columns [ { src_field col path $[0] dest_field name dest_type string }, { src_field col path $[1] dest_field age dest_type int } ] } }那么数据结果表fake1将会像这样nameagecolothera18[a,18]...实现原理ROW 类型的源字段在 JsonPathTransform.java 中被序列化为row.getFields()对应的 JSON 数组因此$[0]、$[1]即对应行内各字段的位置索引。E2E 测试 TestJsonPathTransformIT.java 中的testNestedRow用例即为该场景的端到端验证。配置异常数据处理策略您可以配置row_error_handle_way与column_error_handle_way来处理异常数据两者都是非必填项。row_error_handle_way配置对行数据内所有数据异常进行处理column_error_handle_way配置对某列数据异常进行处理优先级高于row_error_handle_way。从源码看列级策略的判断逻辑位于 JsonPathTransform.java当 JSONPath 执行抛出JsonPathException时若该列配置了column_error_handle_way且允许跳过allowSkip()即SKIP则记录 debug 日志并返回null否则包装为ErrorDataTransformException携带JSON_PATH_COMPILE_ERROR错误码继续抛出由上层MultipleFieldOutputTransform依据行级策略决定是否跳过整行。相关的错误码定义在 JsonPathTransformErrorCode.java如JSONPATH_ERROR_CODE-01columns 不能为空、JSONPATH_ERROR_CODE-03path 不能为空、JSONPATH_ERROR_CODE-05path 无效等。跳过异常数据行配置跳过任意列有异常的整行数据transform { JsonPath { row_error_handle_way SKIP columns [ { src_field json_data path $.f1 dest_field json_data_f1 }, { src_field json_data path $.f2 dest_field json_data_f2 } ] } }跳过部分异常数据列配置仅对json_data_f1列数据异常跳过填充空值其他列数据异常继续抛出异常中断处理程序transform { JsonPath { row_error_handle_way FAIL columns [ { src_field json_data path $.f1 dest_field json_data_f1 column_error_handle_way SKIP }, { src_field json_data path $.f2 dest_field json_data_f2 } ] } }部分列异常跳过整行配置仅对json_data_f1列数据异常跳过整行数据其他列数据异常继续抛出异常中断处理程序transform { JsonPath { row_error_handle_way FAIL columns [ { src_field json_data path $.f1 dest_field json_data_f1 column_error_handle_way SKIP_ROW }, { src_field json_data path $.f2 dest_field json_data_f2 } ] } }行为验证单元测试 JsonPathTransformTest.java 的testErrorHandleWay覆盖了上述全部组合——行级SKIP时异常行返回null整行被跳过列级SKIP时该列填充null而输出行仍存在列级SKIP_ROW时返回null列级FAIL覆盖行级SKIP时仍抛出异常印证列级优先级更高。E2E 测试中的json_path_with_error_handle_way.conf用例也验证了真实运行环境下的行为。在完整管道中的使用方式JsonPath 插件通常配置在transform块内位于 source 与 sink 之间。一个完整的最小配置骨架如下source { FakeSource { schema { fields { data string } } } } transform { JsonPath { columns [ { src_field data path $.user.name dest_field user_name dest_type string } ] } } sink { Console {} }运行seatunnel命令提交任务后即可在结果表中看到新增的user_name列。插件在启动时会先完成三项初始化见 JsonPathTransform.java定位源字段索引、构建输出列类型、创建类型转换器同时工厂层的ColumnsValidatorJsonPathTransformFactory.java会在配置校验阶段拦截path、src_field、dest_field为空的非法配置保证作业尽早失败而不是运行到一半才报错。更新日志添加 JsonPath 转换含基础类型提取、SeatunnelRow 提取、批量字段提取、异常数据策略等能力。相关资源插件源码配置解析与校验插件工厂与选项规则单元测试端到端测试转换插件通用选项赞分享数据集成ETL大数据批处理流处理变更数据捕获【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/GitHub_Trending/se/seatunnel点击查看免费下载相关推荐终极指南SeaTunnel中的JsonPath数据转换插件详解终极指南SeaTunnel中的JsonPath数据转换插件详解 SeaTunnel作为一款强大的数据集成平台其JsonPath数据转换插件提供了简单高效的J数据集成ETL大数据批处理流处理变更数据捕获Apache SeaTunnel中的JSONPath转换插件详解Apache SeaTunnel中的JSONPath转换插件详解 什么是JSONPath转换插件 JSONPath转换插件是Apache SeaTunnel数据数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel RegexExtract 转换插件实战用正则捕获组从字段中提取多列数据SeaTunnel RegexExtract 转换插件实战用正则捕获组从字段中提取多列数据 RegexExtract transform plugin 导读数据集成ETL大数据批处理流处理变更数据捕获上一篇完全免费永久保存微信聊天记录的终极解决方案WeChatMsg完整指南下一篇NetExec代码质量测量与提升创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考