Flink DataStream JSON 格式指南:JsonSerializationSchema 与 JsonDeserializationSchema 实战
2026/9/20 17:31:20 网站建设 项目流程
  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】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> <groupId>org.apache.flink</groupId> <artifactId>flink-json</artifactId> <version>{{< version >}}</version> <scope>provided</scope> </dependency>
  • <scope>provided</scope>表示该依赖在编译期使用、运行期由 Flink 发行包提供,避免与集群自带的版本冲突;若使用 IDE 本地调试或构建 fat-jar,可临时去掉该 scope。
  • 当前仓库中该模块的版本为2.0-SNAPSHOT,其artifactId与依赖关系可参见 flink-formats/flink-json/pom.xml。从该文件可以看到flink-json内部依赖flink-shaded-jackson(Jackson 的 Shaded 版本,防止与用户自身引入的 Jackson 冲突),并可选地关联flink-table-commonflink-connector-files(用于 Table API 与文件系统格式工厂)。
  • 对于 PyFlink 用户,无需额外安装任何包即可直接使用 JSON 格式能力,相关 Python API 定义在 flink-python/pyflink/datastream/formats/json.py。

核心原理:JsonSerializationSchema 与 JsonDeserializationSchema

Flink 通过JsonSerializationSchema/JsonDeserializationSchema支持 JSON 记录的读写。这两个类底层依赖 Jackson 库,能够处理 Jackson 支持的一切类型,包括但不限于POJOObjectNode

反序列化:JsonDeserializationSchema

JsonDeserializationSchema实现了AbstractDeserializationSchema<T>,其核心逻辑非常简洁(见 JsonDeserializationSchema.java):

@Override public T deserialize(byte[] message) throws IOException { return mapper.readValue(message, clazz); }

也就是说,它把每条消息的byte[]直接交给 Jackson 的ObjectMapper.readValue转换成目标类实例。该类提供两类构造函数:

  • JsonDeserializationSchema(Class<T> clazz):按类反序列化,例如 POJO;
  • JsonDeserializationSchema(TypeInformation<T> typeInformation):按TypeInformation反序列化。

JsonDeserializationSchema可用于任何支持DeserializationSchema的连接器。例如与KafkaSource配合,将 Kafka 中的 JSON 消息反序列化为SomePojo

JsonDeserializationSchema<SomePojo> jsonFormat = new JsonDeserializationSchema<>(SomePojo.class); KafkaSource<SomePojo> source = KafkaSource.<SomePojo>builder() .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)

序列化:JsonSerializationSchema

JsonSerializationSchema实现了SerializationSchema<T>(见 JsonSerializationSchema.java),序列化时同样委托给 Jackson:

@Override 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 消息写入 Kafka:

JsonSerializationSchema<SomePojo> jsonFormat = new JsonSerializationSchema<>(); KafkaSink<SomePojo> source = KafkaSink.<SomePojo>builder() .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 都提供接收SerializableSupplier<ObjectMapper>的构造函数,该参数充当 ObjectMapper 的工厂。借助它你可以对创建的 mapper 拥有完全控制权:启用 / 禁用各类 Jackson 特性,或注册模块以扩展支持的类型、增加额外功能。

例如,按 key 排序输出 JSON 字段,并注册ParameterNamesModule以便 POJO 使用构造器参数名完成反序列化绑定:

JsonSerializationSchema<SomeClass> jsonFormat = new JsonSerializationSchema<>( () -> new ObjectMapper() .enable(SerializationFeature.ORDER_MAP_ENTRIES_BY_KEYS) .registerModule(new ParameterNamesModule()));

同样的方式也适用于JsonDeserializationSchema的构造器:

JsonDeserializationSchema<SomeClass> 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类型或参数名绑定等能力的支持。

注意:SerializableSupplier<ObjectMapper>是一个可序列化的函数式接口,它返回的 mapper 工厂会随算子一同分发到各 TaskManager,因此 lambda 内部捕获的对象也必须可序列化。

深入源码:JsonRowSerializationSchema / JsonRowDeserializationSchema(Row 类型)

除了面向 POJO 的通用 Schema,flink-json模块还提供面向 FlinkRow类型的JsonRowSerializationSchemaJsonRowDeserializationSchema。它们通过内部“运行时转换器”(Runtime Converter)在Row与 JacksonJsonNode之间做映射:

  • 序列化时,createConverter按字段类型将Row逐字段转换为ObjectNode,支持INTLONGDOUBLEFLOATSHORTBYTESTRINGBOOLEANBIG_DECBIG_INT以及各类时间类型(SQL_DATESQL_TIMESQL_TIMESTAMPLOCAL_DATELOCAL_TIMELOCAL_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'。它实现了DeserializationFormatFactorySerializationFormatFactory两个接口,并在内部选择性能更优的JsonParserRowDataDeserializationSchema(基于 JacksonJsonParser的流式解析)或传统的JsonRowDataDeserializationSchema

格式相关的全部可选参数定义在 JsonFormatOptions.java,这些参数同样可以为你理解 DataStream 场景下的 JSON 行为提供参考:

参数名类型默认值说明
fail-on-missing-fieldBooleanfalse是否在字段缺失时解析失败;false时缺失字段置为 null
ignore-parse-errorsBooleanfalse是否跳过解析错误的字段/行而不是抛异常;为true时出错字段置为 null
map-null-key.modeStringFAILMap 数据遇到 null key 时的处理:FAIL抛异常 /DROP丢弃该条目 /LITERAL用字面量替换
map-null-key.literalString"null"map-null-key.modeLITERAL时使用的 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 的构造函数可见,ignoreParseErrorsfailOnMissingField同时为true时会被直接判定为非法配置并抛出IllegalArgumentException,这是因为两者语义互斥——一个要求“出错即失败”,另一个要求“出错即忽略”。

PyFlink:使用 JSON Row 格式与 Kafka 集成

在 PyFlink 中,JsonRowSerializationSchemaJsonRowDeserializationSchema内建支持Row类型,对应实现位于 flink-python/pyflink/datastream/formats/json.py。两者均通过 Builder 模式构建,底层调用 Java 侧org.apache.flink.formats.json同名类。

在 KafkaSource 中反序列化 JSON 为 Row

row_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.convert);
  • fail_on_missing_field():字段缺失时解析失败;
  • ignore_parse_errors():解析失败时不抛异常。

注意:Python 侧fail_on_missing_fieldignore_parse_errors若同时开启,最终会经由 Java 侧校验抛错,因此二选一即可。

在 KafkaSink 中将 Row 序列化为 JSON

row_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 JsonDeserializationSchema<SomePojo> jsonFormat = new JsonDeserializationSchema<>(SomePojo.class); KafkaSource<SomePojo> source = KafkaSource.<SomePojo>builder() .setBootstrapServers("localhost:9092") .setTopics("input-topic") .setGroupId("json-demo") .setValueOnlyDeserializer(jsonFormat) .build(); // 3. 处理数据(此处仅打印) DataStream<SomePojo> stream = env.fromSource( source, WatermarkStrategy.noWatermarks(), "kafka-json-source"); stream.map(pojo -> "name=" + pojo.name + ", age=" + pojo.age).print(); // 4. 构造序列化 Schema 并接入 KafkaSink JsonSerializationSchema<SomePojo> outFormat = new JsonSerializationSchema<>(); KafkaSink<SomePojo> sink = KafkaSink.<SomePojo>builder() .setBootstrapServers("localhost:9092") .setRecordSerializer( new KafkaRecordSerializationSchemaBuilder<SomePojo>() .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 等连接器直接组合,代码量最少;
  • 进阶控制:通过SerializableSupplier<ObjectMapper>注入自定义 ObjectMapper,可开关 Jackson 特性、注册模块,满足排序输出、时间类型、参数名绑定等定制需求;
  • Row 场景:PyFlink 内建JsonRowDeserializationSchema/JsonRowSerializationSchema(Builder 模式)开箱即用;Java DataStream 中同名 Row Schema 已标记废弃,建议改用 Table API 或自定义 Schema;
  • 生产配置:字段缺失、解析错误、map 的 null key、时间戳格式等行为可通过格式参数精细控制,其默认值与语义以 JsonFormatOptions.java 为准,并用单元测试(如 JsonSerDeSchemaTest.java、JsonRowDataSerDeSchemaTest.java)验证序列化、反序列化与往返一致性。
  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/flink
点击查看免费下载

相关推荐

上一篇:LaTeX Workshop终极指南:在VS Code中实现高效专业排版的完整方案
下一篇:YTPro的JavaScript接口:原生功能如何通过JS调用

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询