Flink DataStream JSON 格式指南:JsonSerializationSchema 与 JsonDeserializationSchema 实战

发布时间:2026/9/20 17:31:27

Flink DataStream JSON 格式指南:JsonSerializationSchema 与 JsonDeserializationSchema 实战 大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载导读本文聚焦 Apache Flink DataStream API 中最常用的数据交换格式——JSON完整讲解flink-json模块提供的JsonSerializationSchema/JsonDeserializationSchema的使用方法、底层 Jackson 机制、自定义 ObjectMapper 的进阶技巧以及 PyFlink 下JsonRowSerializationSchema/JsonRowDeserializationSchema的用法。读完本文你将能够在 Kafka、FileSystem 等任意支持序列化/反序列化协议的连接器上用最少的代码完成 POJO 与 JSON 字节流的互转并掌握字段缺失、解析失败、时间戳格式等生产级细节的配置方法。添加依赖要在 Java / Scala 项目中使用 JSON 格式需要在工程的pom.xml中引入flink-json依赖dependency groupIdorg.apache.flink/groupId artifactIdflink-json/artifactId version{{ version }}/version scopeprovided/scope /dependencyscopeprovided/scope表示该依赖在编译期使用、运行期由 Flink 发行包提供避免与集群自带的版本冲突若使用 IDE 本地调试或构建 fat-jar可临时去掉该 scope。当前仓库中该模块的版本为2.0-SNAPSHOT其artifactId与依赖关系可参见 flink-formats/flink-json/pom.xml。从该文件可以看到flink-json内部依赖flink-shaded-jacksonJackson 的 Shaded 版本防止与用户自身引入的 Jackson 冲突并可选地关联flink-table-common与flink-connector-files用于 Table API 与文件系统格式工厂。对于 PyFlink 用户无需额外安装任何包即可直接使用 JSON 格式能力相关 Python API 定义在 flink-python/pyflink/datastream/formats/json.py。核心原理JsonSerializationSchema 与 JsonDeserializationSchemaFlink 通过JsonSerializationSchema/JsonDeserializationSchema支持 JSON 记录的读写。这两个类底层依赖 Jackson 库能够处理 Jackson 支持的一切类型包括但不限于POJO和ObjectNode。反序列化JsonDeserializationSchemaJsonDeserializationSchema实现了AbstractDeserializationSchemaT其核心逻辑非常简洁见 JsonDeserializationSchema.javaOverride public T deserialize(byte[] message) throws IOException { return mapper.readValue(message, clazz); }也就是说它把每条消息的byte[]直接交给 Jackson 的ObjectMapper.readValue转换成目标类实例。该类提供两类构造函数JsonDeserializationSchema(ClassT clazz)按类反序列化例如 POJOJsonDeserializationSchema(TypeInformationT typeInformation)按TypeInformation反序列化。JsonDeserializationSchema可用于任何支持DeserializationSchema的连接器。例如与KafkaSource配合将 Kafka 中的 JSON 消息反序列化为SomePojoJsonDeserializationSchemaSomePojo jsonFormat new JsonDeserializationSchema(SomePojo.class); KafkaSourceSomePojo source KafkaSource.SomePojobuilder() .setValueOnlyDeserializer(jsonFormat) ...要点POJO 必须提供无参构造函数且字段要有对应的 getter / setter否则 Jackson 无法完成绑定JSON 中多余字段默认会被忽略POJO 中未出现在 JSON 里的字段默认保持为 null不报错若只需读取 JSON 树而不关心类型绑定可以反序列化为ObjectNode。仓库中提供JsonNodeDeserializationSchema见 JsonNodeDeserializationSchema.java它等价于new JsonDeserializationSchema(ObjectNode.class)之后可通过objectNode.get(name).as(type)访问字段——不过该类的 Javadoc 已明确建议直接使用JsonDeserializationSchema(ObjectNode.class)。序列化JsonSerializationSchemaJsonSerializationSchema实现了SerializationSchemaT见 JsonSerializationSchema.java序列化时同样委托给 JacksonOverride public byte[] serialize(T element) { try { return mapper.writeValueAsBytes(element); } catch (JsonProcessingException e) { throw new RuntimeException( String.format(Could not serialize value %s., element), e); } }它提供了默认无参构造器JsonSerializationSchema()内部使用new ObjectMapper()。JsonSerializationSchema可用于任何支持SerializationSchema的连接器例如与KafkaSink配合把SomePojo序列化为 JSON 消息写入 KafkaJsonSerializationSchemaSomePojo jsonFormat new JsonSerializationSchema(); KafkaSinkSomePojo source KafkaSink.SomePojobuilder() .setRecordSerializer( new KafkaRecordSerializationSchemaBuilder() .setValueSerializationSchema(jsonFormat) ...生命周期与可序列化性两个 Schema 都实现了open(InitializationContext context)方法在其中通过工厂mapperFactory.get()创建ObjectMapper字段用transient修饰不参与算子状态序列化。这意味着自定义的ObjectMapper配置在算子open()时才生效与算子并行度、状态恢复等机制完全兼容。仓库中的单元测试 JsonSerDeSchemaTest.java 展示了标准用法先open(new DummyInitializationContext())再调用serialize/deserialize并验证了{x:34,y:hello}这一 POJO 序列化、反序列化及往返round-trip的一致性。自定义 Mapper精细化控制 JSON 行为两个 Schema 都提供接收SerializableSupplierObjectMapper的构造函数该参数充当 ObjectMapper 的工厂。借助它你可以对创建的 mapper 拥有完全控制权启用 / 禁用各类 Jackson 特性或注册模块以扩展支持的类型、增加额外功能。例如按 key 排序输出 JSON 字段并注册ParameterNamesModule以便 POJO 使用构造器参数名完成反序列化绑定JsonSerializationSchemaSomeClass jsonFormat new JsonSerializationSchema( () - new ObjectMapper() .enable(SerializationFeature.ORDER_MAP_ENTRIES_BY_KEYS) .registerModule(new ParameterNamesModule()));同样的方式也适用于JsonDeserializationSchema的构造器JsonDeserializationSchemaSomeClass jsonFormat new JsonDeserializationSchema( SomeClass.class, () - new ObjectMapper() .disable(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES) .registerModule(new JavaTimeModule()));常见自定义场景SerializationFeature.ORDER_MAP_ENTRIES_BY_KEYS序列化时按键排序便于生成确定性输出、辅助比对DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES控制遇到未知字段是否抛异常默认忽略注册JavaTimeModule/ParameterNamesModule等模块扩展对java.time类型或参数名绑定等能力的支持。注意SerializableSupplierObjectMapper是一个可序列化的函数式接口它返回的 mapper 工厂会随算子一同分发到各 TaskManager因此 lambda 内部捕获的对象也必须可序列化。深入源码JsonRowSerializationSchema / JsonRowDeserializationSchemaRow 类型除了面向 POJO 的通用 Schemaflink-json模块还提供面向 FlinkRow类型的JsonRowSerializationSchema与JsonRowDeserializationSchema。它们通过内部“运行时转换器”Runtime Converter在Row与 JacksonJsonNode之间做映射序列化时createConverter按字段类型将Row逐字段转换为ObjectNode支持INT、LONG、DOUBLE、FLOAT、SHORT、BYTE、STRING、BOOLEAN、BIG_DEC、BIG_INT以及各类时间类型SQL_DATE、SQL_TIME、SQL_TIMESTAMP、LOCAL_DATE、LOCAL_TIME、LOCAL_DATE_TIME、嵌套Row、对象数组、原始byte[]等见 JsonRowSerializationSchema.java时间类型默认按 ISO-8601 / RFC3339 风格的字符串输出如LocalDateTime使用 RFC3339 时间戳格式LocalTime使用 RFC3339 时间格式未通过 JSON Schema 显式描述的类型如 POJO会走mapper.valueToTree(object)的 fallback 转换。从源码的Deprecated注解及 Javadoc 可知这两个 Row Schema 最初是为 Table API 用户开发的官方已声明不再为 DataStream API 用户维护。DataStream 场景建议要么使用 Table API要么自行实现SerializationSchema/DeserializationSchema例如直接基于本文的通用 Schema 或自定义 Jackson 逻辑。Table API / SQL 侧的 JSON 格式与可选参数flink-json同时以格式工厂Format Factory的形式深度集成 Table API / SQL。入口为 JsonFormatFactory.java标识符为json例如 Kafka DDL 中format json。它实现了DeserializationFormatFactory与SerializationFormatFactory两个接口并在内部选择性能更优的JsonParserRowDataDeserializationSchema基于 JacksonJsonParser的流式解析或传统的JsonRowDataDeserializationSchema。格式相关的全部可选参数定义在 JsonFormatOptions.java这些参数同样可以为你理解 DataStream 场景下的 JSON 行为提供参考参数名类型默认值说明fail-on-missing-fieldBooleanfalse是否在字段缺失时解析失败false时缺失字段置为 nullignore-parse-errorsBooleanfalse是否跳过解析错误的字段/行而不是抛异常为true时出错字段置为 nullmap-null-key.modeStringFAILMap 数据遇到 null key 时的处理FAIL抛异常 /DROP丢弃该条目 /LITERAL用字面量替换map-null-key.literalStringnullmap-null-key.mode为LITERAL时使用的 key 字面量timestamp-format.standardStringSQL时间戳格式SQLyyyy-MM-dd HH:mm:ss.s{precision}或ISO-8601yyyy-MM-ddTHH:mm:ss.s{precision}encode.decimal-as-plain-numberBooleanfalse是否把所有 decimal 编码为普通数字而非可能的科学计数法encode.ignore-null-fieldsBooleanfalse编码时是否忽略 null 字段decode.json-parser.enabledBooleantrue是否使用 JacksonJsonParser以更高性能解码 JSON从 JsonRowDeserializationSchema.java 的构造函数可见ignoreParseErrors与failOnMissingField同时为true时会被直接判定为非法配置并抛出IllegalArgumentException这是因为两者语义互斥——一个要求“出错即失败”另一个要求“出错即忽略”。PyFlink使用 JSON Row 格式与 Kafka 集成在 PyFlink 中JsonRowSerializationSchema和JsonRowDeserializationSchema内建支持Row类型对应实现位于 flink-python/pyflink/datastream/formats/json.py。两者均通过 Builder 模式构建底层调用 Java 侧org.apache.flink.formats.json同名类。在 KafkaSource 中反序列化 JSON 为 Rowrow_type_info Types.ROW_NAMED([name, age], [Types.STRING(), Types.INT()]) json_format JsonRowDeserializationSchema.builder().type_info(row_type_info).build() source KafkaSource.builder() \ .set_value_only_deserializer(json_format) \ .build()type_info用于声明结果的Row结构其字段名将用于匹配 JSON 属性名。PyFlink 的 Builder 还支持链式调用json_schema(json_schema: str)基于 JSON Schema 声明结果类型内部调用 Java 侧JsonRowSchemaConverter.convertfail_on_missing_field()字段缺失时解析失败ignore_parse_errors()解析失败时不抛异常。注意Python 侧fail_on_missing_field与ignore_parse_errors若同时开启最终会经由 Java 侧校验抛错因此二选一即可。在 KafkaSink 中将 Row 序列化为 JSONrow_type_info Types.ROW_NAMED([name, age], [Types.STRING(), Types.INT()]) json_format JsonRowSerializationSchema.builder().with_type_info(row_type_info).build() sink KafkaSink.builder() \ .set_record_serializer( KafkaRecordSerializationSchema.builder() .set_topic(test) .set_value_serialization_schema(json_format) .build() ) \ .build()序列化后的byte[]消息可以由JsonRowDeserializationSchema反向解析二者配合即可实现 Row 与 JSON 的无损互转。完整可运行示例Kafka JSON 读写链路将上述 Java 片段整合为一个完整的端到端流程POJO 定义省略 getter / setter// 1. 定义 POJO public static class SomePojo { public String name; public int age; public SomePojo() {} // Jackson 需要无参构造器 } // 2. 构造反序列化 Schema 并接入 KafkaSource JsonDeserializationSchemaSomePojo jsonFormat new JsonDeserializationSchema(SomePojo.class); KafkaSourceSomePojo source KafkaSource.SomePojobuilder() .setBootstrapServers(localhost:9092) .setTopics(input-topic) .setGroupId(json-demo) .setValueOnlyDeserializer(jsonFormat) .build(); // 3. 处理数据此处仅打印 DataStreamSomePojo stream env.fromSource( source, WatermarkStrategy.noWatermarks(), kafka-json-source); stream.map(pojo - name pojo.name , age pojo.age).print(); // 4. 构造序列化 Schema 并接入 KafkaSink JsonSerializationSchemaSomePojo outFormat new JsonSerializationSchema(); KafkaSinkSomePojo sink KafkaSink.SomePojobuilder() .setBootstrapServers(localhost:9092) .setRecordSerializer( new KafkaRecordSerializationSchemaBuilderSomePojo() .setTopic(output-topic) .setValueSerializationSchema(outFormat) .build()) .build(); stream.sinkTo(sink); env.execute(flink-json-datastream-demo);这条链路验证了文档所述的两个核心事实JsonDeserializationSchema适配任何支持DeserializationSchema的连接器JsonSerializationSchema适配任何支持SerializationSchema的连接器二者配合即可完成“Kafka 读 JSON → 处理 → 写 JSON”的典型数据管道。小结POJO 场景优先使用JsonDeserializationSchema(SomePojo.class)与JsonSerializationSchema()与 Kafka / FileSystem 等连接器直接组合代码量最少进阶控制通过SerializableSupplierObjectMapper注入自定义 ObjectMapper可开关 Jackson 特性、注册模块满足排序输出、时间类型、参数名绑定等定制需求Row 场景PyFlink 内建JsonRowDeserializationSchema/JsonRowSerializationSchemaBuilder 模式开箱即用Java DataStream 中同名 Row Schema 已标记废弃建议改用 Table API 或自定义 Schema生产配置字段缺失、解析错误、map 的 null key、时间戳格式等行为可通过格式参数精细控制其默认值与语义以 JsonFormatOptions.java 为准并用单元测试如 JsonSerDeSchemaTest.java、JsonRowDataSerDeSchemaTest.java验证序列化、反序列化与往返一致性。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Apache Flink DataStream Formats 全解析Avro、Parquet、Text、Hadoop 与 Azure Table 编码格式实战指南Apache Flink DataStream Formats 全解析Avro、Parquet、Text、Hadoop 与 Azure Table 编码格式实大数据流处理批处理数据工程Apache Flink DataStream Parquet 格式全解析RowData 向量化读取与 Avro 记录读取实战Apache Flink DataStream Parquet 格式全解析RowData 向量化读取与 Avro 记录读取实战 本指南围绕 Apache Fl大数据流处理批处理数据工程Apache Flink CDC PostgreSQL Connector 全指南从建表配置到增量快照与 DataStream 实战Apache Flink CDC PostgreSQL Connector 全指南从建表配置到增量快照与 DataStream 实战 PostgreSQL C后端数据集成大数据流处理变更数据捕获数据同步上一篇LaTeX Workshop终极指南在VS Code中实现高效专业排版的完整方案下一篇YTPro的JavaScript接口原生功能如何通过JS调用创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
延伸阅读

更多相关文章

2026/9/20 17:26:26

radare2 libr 架构解析:一张图看懂核心库 API 依赖关系

逆向工程网络安全 【免费下载链接】radare2 UNIX-like reverse engineering framework and command-line toolset 项目地址: https://gitcode.com/gh_mirrors/ra/radare2 点击查看 免费下载 导读:radare2 的功能被拆分为多个以 r_ 前缀命名的核心库&…

2026/9/20 17:26:26

Linux驱动自动加载全解析:从内核模块到设备树实战

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/20 18:26:33

Bili.Uwp 上手指南:Windows 上跑起自己的 UWP 哔哩哔哩客户端

Bili.Uwp 上手指南:Windows 上跑起自己的 UWP 哔哩哔哩客户端 【免费下载链接】Bili.Uwp 适用于新系统UI的哔哩 项目地址: https://gitcode.com/GitHub_Trending/bi/Bili.Uwp Bili.Uwp(仓库内名为“哔哩”)是一款用 C# 和 UWP 框架开发…

2026/9/20 18:26:33

Atlas 300V实战:从ONNX到OM,用CANN部署YOLO推理模型

1. 先别急着部署,Atlas 300V到底是个什么"卡"看到"atlas 300v 24g 是运算加速卡吗"这个问题的时候,我基本能猜到提问的人正处于哪个阶段:手里刚刚拿到一块Atlas 300V,插到服务器上,正准备像装NVID…

2026/9/20 18:26:33

KataGo围棋AI配置与优化全指南

1. 项目概述KataGo作为当前最强大的开源围棋AI之一,其神经网络架构和搜索算法在棋力表现上已经超越了许多商业软件。不同于传统围棋引擎,KataGo采用蒙特卡洛树搜索(MCTS)与深度神经网络结合的架构,支持自定义规则和让子…

2026/9/20 18:26:33

AIGC模型部署方案详解:本地、云端与混合架构的选型与实战

过去一年里,找我咨询“AIGC模型部署”的人比预想中多得多。多数人不是不会跑代码,而是卡在第一个选择题上:到底是买一台机器在本地部署,还是直接调云端API,又或者两边各放一部分形成混合架构。这个决策直接影响后面的成…

2026/9/20 18:21:33

Edge垂直标签页设置教程:宽屏效率提升与标签管理技巧

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/20 0:04:49

GAMP 5 基于风险的计算机化系统验证:软件分类与审计追踪实践

简介:《A Risk-Based Approach to Compliant GxP Computerized Systems》即业内熟知的GAMP 5指南,面向制药企业质量与IT合规人员、验证工程师及计算机化系统管理者,用于解决GxP法规环境下系统合规性难以科学落地的问题。文档以风险管理为主线…

2026/9/20 0:04:49

安全托管MSSP实战:从静态防御到人机协同的攻防运营与应急响应

简介:这份PPT围绕互联网业务安全托管服务展开,面向企业安全负责人、IT运维人员及关注MSSP/MSS选型的读者,重点回应传统安全过度依赖人工、碎片化静态防御难以对抗产业化攻击等痛点。资源共1个pptx文件,包体约30.63MB,以…

2026/9/20 0:04:49

GAMP 5 基于风险的计算机化系统验证:软件分类与审计追踪实践

简介:《A Risk-Based Approach to Compliant GxP Computerized Systems》即业内熟知的GAMP 5指南,面向制药企业质量与IT合规人员、验证工程师及计算机化系统管理者,用于解决GxP法规环境下系统合规性难以科学落地的问题。文档以风险管理为主线…

2026/9/20 0:04:49

安全托管MSSP实战:从静态防御到人机协同的攻防运营与应急响应

简介:这份PPT围绕互联网业务安全托管服务展开,面向企业安全负责人、IT运维人员及关注MSSP/MSS选型的读者,重点回应传统安全过度依赖人工、碎片化静态防御难以对抗产业化攻击等痛点。资源共1个pptx文件,包体约30.63MB,以…

2026/9/20 4:54:47

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行,Type-C接口算是典型的“看着简单,做起来全坑”的东西。光引脚就24个,高低速信号、电源、控制线全部塞在一个小小的连接器里,如果PCB布局不做规划,打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/20 5:01:23

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/20 5:09:33

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
咨询二维码