大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载导读本文以 flink-python/docs/reference/pyflink.common/serializer.rst 所定义的 API 清单为主线系统梳理 PyFlink 中面向运行时状态管理的TypeSerializer体系以及面向连接器数据交换的SerializationSchema/DeserializationSchema/Encoder/BulkWriterFactory体系。你将掌握两类序列化接口的职责边界与 Python 源码实现、SimpleStringSchema的编码原理与字符集定制方法以及在 Kafka 与文件 Sink 中组合使用它们的完整实战方式。一、PyFlink 的两套序列化接口分清职责边界PyFlink 的序列化能力分布在两个不同的 Python 模块中承担两种截然不同的任务模块核心类职责pyflink.common.serializerTypeSerializer、VoidNamespaceSerializer供 Flink 运行时内部处理数据类型状态、命名空间等的序列化/反序列化pyflink.common.serializationSerializationSchema、DeserializationSchema、SimpleStringSchema、Encoder、BulkWriterFactory、RowDataBulkWriterFactory供外部连接器Kafka、文件系统等把数据对象转换为字节流或把字节流还原为数据对象这两个模块都在 pyflink/common/init.py 中被统一导出到pyflink.common命名空间下例如from pyflink.common import SimpleStringSchema, TypeSerializer因此用户代码中两者可以混合使用但语义上运行时内部与连接器边界的分工务必分清。二、TypeSerializer运行时数据类型的序列化契约2.1 抽象基类与两个抽象方法TypeSerializer定义在 pyflink/common/serializer.py是一个继承ABC并携带泛型参数T的抽象基类。其源码 docstring 明确说明该接口描述了一个数据类型要被 Flink 运行时处理所必需的方法具体就是序列化与反序列化两个方法from abc import abstractmethod, ABC from io import BytesIO from typing import TypeVar, Generic T TypeVar(T) class TypeSerializer(ABC, Generic[T]): abstractmethod def serialize(self, element: T, stream: BytesIO) - None: Serializes an element to the output stream. pass abstractmethod def deserialize(self, stream: BytesIO) - T: Returns a deserialized element from the input stream. pass关键设计点面向流式读写serialize把元素写入BytesIO输出流返回Nonedeserialize从输入流读取并返回元素。这与SerializationSchema面向元素 ↔ 字节数组的批量转换不同。参数化类型通过Generic[T]声明操作的元素类型子类实例化时需给出具体类型例如VoidNamespaceSerializer声明为TypeSerializer[bytes]。2.2 内置的相等性、哈希与字符串表示基类还提供了默认的__eq__/__ne__/__repr__/__hash__实现__eq__只有两个实例属于同一类且__dict__完全相等时才相等__repr__返回ClassName()形式的简洁描述__hash__基于str(self)计算哈希与__eq__配套保证序列化器可以作为字典键或集合元素使用。这些实现的意义在于Flink 运行时经常需要比较、缓存、分发序列化器例如按类型注册、状态恢复时校验序列化器一致性Python 侧的默认语义提供了稳定的比较与哈希行为。2.3_get_coder()桥接 Beam Coder 协议基类中的_get_coder()方法serializer.py把当前序列化器适配为类似 Apache Beam Coder 的接口def _get_coder(self): serialize_func self.serialize deserialize_func self.deserialize class CoderAdapter(object): def get_impl(self): return CoderAdapterIml() class CoderAdapterIml(object): def encode_nested(self, element): bytes_io BytesIO() serialize_func(element, bytes_io) return bytes_io.getvalue() def decode_nested(self, bytes_data): bytes_io BytesIO(bytes_data) return deserialize_func(bytes_io) return CoderAdapter()它把流式的serialize/deserialize包装成字节数组的encode_nested/decode_nested从而让自定义序列化器可以在 PyFlink 与 Beam 生态互操作如 Fn 执行相关路径中直接复用。2.4 VoidNamespaceSerializer空命名空间的哨兵序列化器VoidNamespaceSerializerserializer.py是TypeSerializer[bytes]的极简实现void b class VoidNamespaceSerializer(TypeSerializer[bytes]): def serialize(self, element: bytes, stream: BytesIO) - None: pass def deserialize(self, stream: BytesIO) - bytes: return void它在序列化时什么都不写入对应空命名空间没有有效载荷反序列化时固定返回b。这类哨兵序列化器在 Flink 中用于表示不需要携带任何数据的状态命名空间PyFlink 以此保证与 Java 侧VoidNamespaceSerializer语义对齐。三、SerializationSchema 与 DeserializationSchema连接器边界的数据转换3.1 两个基类的语义pyflink/common/serialization.py 定义了连接器侧的基类SerializationSchema描述如何把一个数据对象转换为另一种序列化表示。绝大多数数据 Sink例如 Apache Kafka要求数据以特定格式如字节串交给它。DeserializationSchema描述如何把某些数据源例如 Apache Kafka投递的字节消息还原为 Flink 处理的数据类型此外它还描述产出的类型使 Flink 能够据此创建内部序列化器与数据结构来管理该类型。两者都只是持有对应 Java 对象的薄包装SerializationSchema持有_j_serialization_schemaDeserializationSchema持有_j_deserialization_schema真正的编解码逻辑在 Java 侧实现Python 侧通过 JVM 网关调用。3.2 SimpleStringSchema开箱即用的字符串编解码SimpleStringSchema同时继承SerializationSchema与DeserializationSchema是连接器场景中使用频率最高的开箱即用实现serialization.pyclass SimpleStringSchema(SerializationSchema, DeserializationSchema): def __init__(self, charset: str UTF-8): gate_way get_gateway() j_char_set gate_way.jvm.java.nio.charset.Charset.forName(charset) j_simple_string_serialization_schema gate_way \ .jvm.org.apache.flink.api.common.serialization.SimpleStringSchema(j_char_set) SerializationSchema.__init__( self, j_serialization_schemaj_simple_string_serialization_schema) DeserializationSchema.__init__( self, j_deserialization_schemaj_simple_string_serialization_schema)使用要点默认字符集 UTF-8构造时不传参数即使用UTF-8也可传入任意 Java 支持的标准字符集名如SimpleStringSchema(GBK)。一个实例同时具备两重身份它被同时注册为SerializationSchema与DeserializationSchema同一个实例既可挂在 Source 上做反序列化也可挂在 Sink 上做序列化。3.3 底层 Java 实现的行为细节Python 包装对应的 Java 类是org.apache.flink.api.common.serialization.SimpleStringSchema源码位于 flink-core/src/main/java/org/apache/flink/api/common/serialization/SimpleStringSchema.java。其核心行为可对照验证编解码deserialize(byte[] message)返回new String(message, charset)serialize(String element)返回element.getBytes(charset)产出类型getProducedType()返回BasicTypeInfo.STRING_TYPE_INFO即 Flink 会将消息识别为 String 类型流终止判定isEndOfStream(String)恒返回false即普通字符串消息不会触发数据流结束跨节点传输charset字段被声明为transient并通过自定义writeObject/readObject把字符集名以writeUTF/readUTF形式序列化确保序列化器在作业分发/恢复时字符集不丢失。PyFlink 侧的单元测试 pyflink/common/tests/test_serialization_schemas.py 验证了这一往返行为用 UTF-8 编码字符串后经 Java 侧serialize得到字节再经deserialize还原为原字符串。3.4 在 Kafka 连接器中的典型用法SimpleStringSchema最常见的落地场景就是 Kafka。以 pyflink/datastream/connectors/kafka.py 中KafkaSource的 docstring 示例为证 source KafkaSource \ ... .builder() \ ... .set_bootstrap_servers(MY_BOOTSTRAP_SERVERS) \ ... .set_group_id(MY_GROUP) \ ... .set_topics(TOPIC1, TOPIC2) \ ... .set_value_only_deserializer(SimpleStringSchema()) \ ... .set_starting_offsets(KafkaOffsetsInitializer.earliest()) \ ... .build()在写入侧FlinkKafkaProducer同样接受set_value_serialization_schema(SimpleStringSchema())/set_key_serialization_schema(SimpleStringSchema())来把 String 键值编码为 Kafka 消息字节见 kafka.py 与 kafka.py 的示例。四、Encoder 与文件 Sink 的行式写出Encoderserialization.py是文件 Sink 用于把到达的元素实际写入桶内文件的接口其内置工厂方法为staticmethod def simple_string_encoder(charset_name: str UTF-8) - Encoder: ...默认行为simple_string_encoder()对应 Java 侧SimpleStringEncoder源码见 flink-core/src/main/java/org/apache/flink/api/common/serialization/SimpleStringEncoder.java它会对输入元素调用toString()并写入桶文件元素之间以换行符分隔适用场景行式文本文件输出如每行一条记录配合 FileSink 的分桶与滚动策略使用若要输出二进制/批量格式则应改用下一节的BulkWriterFactory。五、BulkWriterFactory 与 RowDataBulkWriterFactory批量文件写入BulkWriterFactoryserialization.py是 JavaBulkWriter.Factory接口的 Python 包装是所有以批量方式把记录写入文件的数据 Sink 的基础接口RowDataBulkWriterFactoryserialization.py是前者的子类专门接收RowData类型记录构造时额外携带row_type表示 Python 侧的 Row 记录必须先转换为RowData才能写入。两者的落地用法体现在 FileSink 的for_bulk_format工厂方法中pyflink/datastream/connectors/file_system.pystaticmethod def for_bulk_format(base_path: str, writer_factory: BulkWriterFactory) \ - FileSink.BulkFormatBuilder: jvm get_gateway().jvm j_path jvm.org.apache.flink.core.fs.Path(base_path) JFileSink jvm.org.apache.flink.connector.file.sink.FileSink builder FileSink.BulkFormatBuilder( JFileSink.forBulkFormat(j_path, writer_factory.get_java_object()) ) if isinstance(writer_factory, RowDataBulkWriterFactory): return builder._with_row_type(writer_factory.get_row_type()) else: return builder从源码结构可以推断当传入的是RowDataBulkWriterFactory时构建器会额外记录row_type以便后续对 Python Row 与 Java RowData 做转换普通BulkWriterFactory则直接透传。因此 PyFlink 中 Parquet、ORC 等需要按列式/批量编码的格式通常都基于此类工厂构造 FileSink。六、实践路径总结如何选择正确的序列化接口你的需求应使用的类模块自定义运行时数据类型/状态序列化TypeSerializer实现serialize/deserializepyflink.common.serializer空命名空间占位VoidNamespaceSerializerpyflink.common.serializerKafka 等字节消息 Source/Sink 的 String 编解码SimpleStringSchema可自定义字符集pyflink.common.serialization行式文本文件写出toString 换行Encoder.simple_string_encoder()pyflink.common.serializationParquet/ORC 等批量格式文件写出BulkWriterFactory/RowDataBulkWriterFactory配合FileSink.for_bulk_formatpyflink.common.serializationpyflink.datastream.connectors.file_system编写自定义TypeSerializer时建议复用基类自带的_get_coder()桥接能力并保持__eq__/__hash__语义一致接入连接器时优先复用SimpleStringSchema与各格式模块提供的 Row Schema 实现如JsonRowDeserializationSchema、CsvRowDeserializationSchema等均由 pyflink/common/init.py 导出到pyflink.common命名空间再按需扩展。相关示例可继续查阅 flink-python/pyflink/examples/datastream/connectors/ 下的kafka_json_format.py、kafka_csv_format.py等完整程序。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Julia 的 Serialization 模块基于 stdlib/Serialization 的序列化与反序列化完全指南Julia 的 Serialization 模块基于 stdlib/Serialization 的序列化与反序列化完全指南 导读 Serialization编程语言编译器语言运行时标准库JIT编译Flink 托管状态自定义序列化完整指南从 TypeSerializer 到状态 Schema 演进Flink 托管状态自定义序列化完整指南从 TypeSerializer 到状态 Schema 演进 本指南面向需要使用自定义序列化方案来管理 Flink 托大数据流处理批处理数据工程Fuel序列化扩展与Kotlinx Serialization的完美结合指南Fuel序列化扩展与Kotlinx Serialization的完美结合指南 Fuel作为Kotlin/Android平台上最简单的HTTP网络库通过与Ko后端网络上一篇Rufus工具解锁Windows 11安装限制的终极解决方案下一篇ant-design Pagination 更多分页省略页码实战大列表分页与折叠页码机制详解创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
企业数字化 ERP 产品动态
相关推荐
m3u8 视频怎么下载?MediaGo 内置浏览器嗅探与 HLS 分片下载完整教程 音视频桌面应用后端 【免费下载链接】mediago 跨平台视频提取工具:支持流媒体下载、视频下载、m3u8 下载及 B站视频下载,提供 Windows 和 Mac 桌面客户端。Cross-platform video extraction tool: Supports streaming download, video download, m3u8 do… · 2026/9/25 3:36:46
computer use类LLM控制计算机工具的探索:用TaoToken统一Key打通Cline配置 /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 3:36:46
Atlas 300V 24G推理加速卡部署YOLO全攻略,手把手绕过踩坑 后台经常有朋友私信我第一句话就问:“Atlas 300V 24G是运算加速卡吗?能不能跑YOLO?”第二句话往往是:“网上说atlas部署yolo很麻烦,是真的吗?”这两个问题我当年刚拿到这张卡时也反复琢磨过。先说结论&… · 2026/9/25 6:49:16
精益与六西格玛:核心差异与协同应用指南 1. 精益与六西格玛的本质差异在制造业和服务业的质量管理实践中,精益(Lean)和六西格玛(Six Sigma)是两种最常被提及的方法论。虽然它们经常被并列讨论,但两者的核心目标和实施路径存在根本性差异。精益起源… · 2026/9/25 6:49:16
C盘又满了?一文教你修改Windows默认安装路径,彻底告别空间告急 C盘又红了,这句话几乎是我每次帮忙解决电脑问题时的开场白。Win10用户最容易遇到的一种情况是:系统盘明明分了128G甚至256G,软件却老是被默认装进C:\Program Files,Windows商店应用也默认往C盘塞,桌面文件、下载文件、… · 2026/9/25 6:49:16
EndNote完全指南:安装、Word插件、文献库管理与高频故障排查 /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 6:49:16
Go语言for-range与switch深度解析与避坑指南 1. 项目概述作为一名长期奋战在Go语言一线的开发者,我见过太多同事在for-range和switch这两个看似简单的语法结构上栽跟头。这些坑往往在代码评审时才会被发现,有时甚至会导致线上事故。今天我们就来彻底剖析这两个语法结构的核心机制,让你在… · 2026/9/25 6:49:10
希格斯场:从上帝粒子到质量起源,粒子物理标准模型的核心枢纽 在对撞机数据和理论物理之间摸爬滚打多年之后,每次被问到“你觉得希格斯场到底是什么”,我都会停一下。因为这个问题看着基础,但真要把它说透,牵扯到的不仅仅是那个著名的“上帝粒子”,更是一整套现代物理学看待世界的… · 2026/9/25 6:49:10
创维E900V22D刷机全攻略:S905L3SB芯片兼容性解析与救砖实战 /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 1:00:31
MQTT协议原理与Broker服务器搭建实战:从Mosquitto到EMQX /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 1:00:37