- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】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-common与flink-connector-files(用于 Table API 与文件系统格式工厂)。 - 对于 PyFlink 用户,无需额外安装任何包即可直接使用 JSON 格式能力,相关 Python API 定义在 flink-python/pyflink/datastream/formats/json.py。
核心原理:JsonSerializationSchema 与 JsonDeserializationSchema
Flink 通过JsonSerializationSchema/JsonDeserializationSchema支持 JSON 记录的读写。这两个类底层依赖 Jackson 库,能够处理 Jackson 支持的一切类型,包括但不限于POJO和ObjectNode。
反序列化: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类型的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-field | Boolean | false | 是否在字段缺失时解析失败;false时缺失字段置为 null |
ignore-parse-errors | Boolean | false | 是否跳过解析错误的字段/行而不是抛异常;为true时出错字段置为 null |
map-null-key.mode | String | FAIL | Map 数据遇到 null key 时的处理:FAIL抛异常 /DROP丢弃该条目 /LITERAL用字面量替换 |
map-null-key.literal | String | "null" | map-null-key.mode为LITERAL时使用的 key 字面量 |
timestamp-format.standard | String | SQL | 时间戳格式:SQL(yyyy-MM-dd HH:mm:ss.s{precision})或ISO-8601(yyyy-MM-ddTHH:mm:ss.s{precision}) |
encode.decimal-as-plain-number | Boolean | false | 是否把所有 decimal 编码为普通数字而非可能的科学计数法 |
encode.ignore-null-fields | Boolean | false | 编码时是否忽略 null 字段 |
decode.json-parser.enabled | Boolean | true | 是否使用 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 为 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_field与ignore_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
相关推荐
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
后端数据集成大数据流处理变更数据捕获数据同步
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考