新闻详情

新闻详情

首页 / 资讯中心 / 详情

Flink PyFlink StatementSet 完整指南:多 Sink 批量写入、统一优化与结果处理

发布时间:2026/9/25 2:44:53来源:尧图网络
Flink PyFlink StatementSet 完整指南:多 Sink 批量写入、统一优化与结果处理
大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载导读本文系统讲解 PyFlink 中StatementSet语句集合的完整用法它允许你将多条 DML 语句或Table对象集中提交由优化器统一优化后作为一个作业job执行是实现单作业多 Sink 分流写入的标准方案。你将掌握create_statement_set的创建方式、add_insert_sql/add_insert/attach_as_datastream三种添加语句的途径、explain预检执行计划、execute提交作业以及通过TableResult与ResultKind消费执行结果的完整实践链路。本文依据当前仓库 flink-python/docs/reference/pyflink.table/statement_set.rst 展开并结合源码与测试用例逐层印证。一、StatementSet 是什么StatementSet是一个用于批量收集、统一优化、单作业提交的 API 容器它接受由 DML 语句如INSERT INTO ... SELECT ...或Table对象定义的管道pipeline规划器planner会对加入集合的所有语句一起做全局优化然后将它们作为一个作业提交执行。这一设计与逐条执行execute_sql形成鲜明对比逐条执行时每条 DML 各自生成、提交一个独立作业相互之间无法共享优化而使用StatementSet后多条写入可以在同一作业内共享数据扫描、算子与状态从而显著降低集群资源开销、减少作业管理成本并保证多条写入的事务性调度统一。核心特征统一优化集合内所有语句由规划器合并优化可跨语句共享公共子计划单作业提交所有语句最终提交为一个 Flink 作业执行后清空调用execute()后已添加的语句会被清除StatementSet不可重复执行已添加的语句。从源码注释看StatementSet自1.11.0版本引入见 statement_set.pyattach_as_datastream则在1.16.0版本新增statement_set.py。如何创建通过TableEnvironment.create_statement_set()创建实例from pyflink.table import EnvironmentSettings, TableEnvironment t_env TableEnvironment.create(EnvironmentSettings.in_streaming_mode()) statement_set t_env.create_statement_set()其底层实现直接调用 Java 侧的TableEnvironment.createStatementSet()并将 Java 对象句柄包装成 Python 端StatementSet见 table_environment.py。二、向 StatementSet 添加语句三种方式1. add_insert_sql直接添加 SQL DML 语句add_insert_sql(stmt)将一条 INSERT SQL 语句加入集合参数为完整的 SQL 文本返回当前StatementSet实例以便链式调用statement_set.add_insert_sql( INSERT INTO first_sink SELECT * FROM source_table WHERE id 3 )底层对应 Java 的addInsertSql(String)statement_set.py。适合已经以 SQL 形态存在的管道定义可直接拼接表名、条件等动态片段。2. add_insert基于 Table 对象写入add_insert(target_path_or_descriptor, table, overwriteFalse)把由Table对象定义的管道写入某个目标表目标表有两种指定方式方式一注册表路径字符串。传入已注册 sink 表的路径字符串路径解析规则与use_catalog/use_database一致即按照当前 catalog → 当前 database → 表名解析。方式二TableDescriptor表描述符。传入的TableDescriptor会被注册为一个内联inline / 匿名临时目录表再向其写入数据。该方式从1.14.0起支持。使用描述符可以在写语句的同时声明 sink 的 Schema支持用自定义DataType覆盖自动推导出的列在物理列旁增加元数据列metadata columns声明主键primary key。若声明的 Schema 中不含物理/常规列这些列会被自动推导并隐式置于 Schema 声明的最前面。官方 docstring 中的完整示例stmt_set table_env.create_statement_set() source_table table_env.from_path(SourceTable) sink_descriptor TableDescriptor.for_connector(blackhole) \ .schema(Schema.new_builder() .build()) \ .build() stmt_set.add_insert(sink_descriptor, source_table)overwrite 参数布尔值指示写入是否覆盖已有数据默认False。底层实现按目标类型分别调用 Java 的addInsert(String, Table, boolean)或addInsert(TableDescriptor, Table, boolean)statement_set.py。3. attach_as_datastream将管道挂载到 DataStream 环境attach_as_datastream()1.16.0将所有已添加语句作为一个整体优化并转换为 transformation 挂载到底层StreamExecutionEnvironment上随后需要用pyflink.datastream.StreamExecutionEnvironment.execute()来真正执行。该方法调用后同样会清空已添加的语句statement_set.py。它适用于需要将 Table 管道与 DataStream API 程序混合提交同一作业的场景。三、预检执行计划explain在真正提交作业前可用explain(*extra_details)获取集合内所有语句与Table的AST 与执行计划返回字符串。可用于提交前人工审查 SQL 改写与物理计划是否正确检查不同ExplainDetail带来的额外信息。ExplainDetail支持的可选明细explain_detail.py枚举值含义ExplainDetail.ESTIMATED_COST优化器估算的物理节点代价如行数、CPU、IO、网络、内存ExplainDetail.CHANGELOG_MODE物理节点产生的 changelog 模式如[I,UA,D]ExplainDetail.JSON_EXECUTION_PLAN程序的 JSON 格式执行计划ExplainDetail.PLAN_ADVICE物理计划中潜在风险告警与 SQL 优化器调优建议源码实现中explain 固定使用ExplainFormat.TEXT文本格式并把 Python 端的 detail 枚举转换为 Java 数组后调用 JavaStatementSet.explain(...)statement_set.py。测试用例展示了同时传入多个 detail 的用法test_table_environment_api.pystmt_set t_env.create_statement_set() stmt_set.add_insert_sql(insert into sink1 select * from %s where a 100 % source) stmt_set.add_insert_sql(insert into sink2 select * from %s where a 100 % source) actual stmt_set.explain(ExplainDetail.ESTIMATED_COST, ExplainDetail.CHANGELOG_MODE, ExplainDetail.JSON_EXECUTION_PLAN) self.assertIsInstance(actual, str)四、批量执行executeexecute()将集合内所有语句与Table作为一个批执行返回TableResult对象。执行前会先调用_t_env._before_execute()完成环境预检然后调用 Java 侧StatementSet.execute()statement_set.py。重要语义执行该方法后已添加的语句与Table会被清除因此该StatementSet实例不能再次执行同样的语句如需重复提交应重新create_statement_set()并重新添加。五、完整实战示例单作业多 Sink 分流写入仓库自带的可直接运行的示例位于 multi_sink.py完整展示一个作业、两个 sink、SQL Table 两种添加方式from pyflink.table import (EnvironmentSettings, TableEnvironment, DataTypes) from pyflink.table.udf import udf t_env TableEnvironment.create(EnvironmentSettings.in_streaming_mode()) table t_env.from_elements( elements[(1, Hello), (2, World), (3, Flink), (4, PyFlink)], schema[id, data]) # 定义两个 sink 表 t_env.execute_sql( CREATE TABLE first_sink ( id BIGINT, data VARCHAR ) WITH ( connector print ) ) t_env.execute_sql( CREATE TABLE second_sink ( id BIGINT, data VARCHAR ) WITH ( connector print ) ) # 创建 statement set statement_set t_env.create_statement_set() # 通过 SQL 语句将 id 3 的数据写入 first_sink statement_set.add_insert_sql(INSERT INTO first_sink SELECT * FROM %s WHERE id 3 % table) # 通过 Table 对象 UDF 过滤包含 Flink 的数据写入 second_sink udf(result_typeDataTypes.BOOLEAN()) def contains_flink(data): return Flink in data second_table table.where(contains_flink(table.data)) statement_set.add_insert(second_sink, second_table) # 统一提交执行本地 mini cluster 下建议调用 wait() 等待作业结束 statement_set.execute().wait()该示例同时体现两种添加方式的混用add_insert_sql接收 SQL 文本add_insert接收Table对象Table.where(...)得到过滤后的新表。两个 sink 的写入在同一个作业中完成。说明execute().wait()用于本地 mini cluster 场景下等待作业结束若提交到远程集群通常移除.wait()以便异步监控作业状态。六、消费执行结果TableResult 与 ResultKindStatementSet.execute()返回 TableResult它是对语句执行结果的统一表示1.11.0 引入。对 INSERT 这类 DML 操作TableResult中每个 sink 对应一列列名为插入表名列类型为BIGINT行的值代表受影响行数。常用方法get_job_client()返回与已提交 Flink 作业关联的JobClient可获取作业状态、取消作业等。仅 DML / DQL 语句有对应作业DDL、DCL 等返回None。wait(timeout_msNone)等待结果就绪。SELECT 操作等待首行数据可被本地访问INSERT 操作等待作业完成其结果仅一行其他操作立即返回。可传毫秒超时值。get_table_schema()获取结果表的 Schema。DDL/USE/EXPLAIN 的结果列名为resultSTRING 类型SHOW 的列名为对象名如catalog name、database name、table name、view name、function nameDESCRIBE 返回name/type/null/key/computed column/watermark五类列INSERT 每 sink 一列SELECT 返回所选字段名与类型。get_result_kind()返回ResultKind判断结果类型。collect()返回一个可关闭的行迭代器CloseableIterator用于逐行拉取结果数据。print()以表格形式将结果打印到客户端控制台。collect 的正确姿势对 SELECT 操作作业在全部结果被消费前不会结束必须主动关闭迭代器以取消作业、释放资源对 DML 操作Flink 暂不支持获取真实受影响行数因此每行受影响行数恒为-1unknown且要等作业结束才返回对 DDL/SHOW 等操作不会提交作业get_job_client()为空结果为有界数据。推荐用with语句确保关闭table_result t_env.execute(select ...) with table_result.collect() as results: for result in results: print(result)注意collect()与print()不能在同一个TableResult实例上同时调用二者任选其一获取本地结果。ResultKind结果的两种类型ResultKind定义了执行结果的类型result_kind.pySUCCESS值为 0语句如 DDL、USE执行成功结果仅包含一个简单的OKSUCCESS_WITH_CONTENT值为 1语句如 DML、DQL、SHOW执行成功结果包含重要内容。底层通过 Java 侧org.apache.flink.table.api.ResultKind的枚举对应转换遇到未知 Java 枚举值会抛出异常。get_result_kind()返回该枚举可用于区分纯确认型与带内容型结果。七、源码级补充执行链路与测试印证Python → Java 的桥接结构StatementSet的 Python 实现本质是 Java API 的薄封装实例持有_j_statement_setJavaorg.apache.flink.table.api.StatementSet的 py4j 引用与所属TableEnvironmentadd_insert_sql、add_insert、explain、execute、attach_as_datastream分别映射到 Java 同名方法。这一设计让 PyFlink 与 Flink Table APIJava保持语义一致。测试用例印证test_table_environment_api.py验证两条add_insert_sql后explain返回字符串覆盖多语句 explain 场景test_table_descriptor.py验证add_insert(sink_descriptor, table)execute().wait()的完整链路印证 TableDescriptor 方式的实际可运行性。适用边界提示add_insert的 TableDescriptor 方式需 Flink 1.14.0attach_as_datastream需 1.16.0本文描述以当前仓库代码为准。StatementSet面向多 DML 合并执行场景若仅需单条语句直接使用execute_sql更简洁。结语StatementSet是 PyFlink 面向多 Sink 场景的核心 API通过create_statement_set创建、add_insert_sql/add_insert/attach_as_datastream添加管道、explain预检、execute单作业提交再配合TableResult的wait/collect/print与ResultKind结果类型判定即可实现高效、可观测的批量写入作业。建议在动手实践前阅读仓库中的 multi_sink.py 示例并结合 statement_set.py 与 table_result.py 源码深入理解其行为。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Flink CDC批量写入优化终极指南Sink端性能调优完全教程Flink CDC批量写入优化终极指南Sink端性能调优完全教程 Flink CDC作为流式数据集成工具在实时数据同步场景中表现出色。而Sink端的批量写入后端数据集成大数据流处理变更数据捕获数据同步SeaTunnel Cassandra Sink 实战批量写入、异步批处理与 CDC 入仓的完整配置指南SeaTunnel Cassandra Sink 实战批量写入、异步批处理与 CDC 入仓的完整配置指南 本文基于 SeaTunnel 官方文档 docs/z数据集成ETL大数据批处理流处理变更数据捕获Flink PyFlink 向量化 Python UDFVectorized UDF实战指南从 pandas.Series 批处理到性能调优Flink PyFlink 向量化 Python UDFVectorized UDF实战指南从 pandas.Series 批处理到性能调优 本文围绕 A大数据流处理批处理数据工程上一篇Rcpp未来展望C20特性在R中的集成可能性分析下一篇GitHub_Trending/le/LeetCode-Book训练计划系列问题解题思路创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
网站建设高端定制企业官网
RELATED

相关资讯

更多精彩内容,欢迎继续阅读

较早相关资讯

最新相关资讯

初等矩阵逆矩阵怎么记?交换、倍乘、倍加三种规则一眼看出 2026/9/25 3:23:47

初等矩阵逆矩阵怎么记?交换、倍乘、倍加三种规则一眼看出

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

阅读更多 →
解码prism4cj核心标记化算法:greedy贪婪匹配与lookbehind如何实现精准代码标记 2026/9/25 3:23:47

解码prism4cj核心标记化算法:greedy贪婪匹配与lookbehind如何实现精准代码标记

解码prism4cj核心标记化算法:greedy贪婪匹配与lookbehind如何实现精准代码标记 【免费下载链接】prism4cj 一个轻量的语法高亮库 项目地址: https://gitcode.com/Cangjie-TPC/prism4cj prism4cj 是一个用仓颉语言编写的轻量级语法高亮库,它的核心…

阅读更多 →
BullMQ Elixir 性能基准测试全解析:吞吐量数据、架构原理与实战调优指南 2026/9/25 3:23:47

BullMQ Elixir 性能基准测试全解析:吞吐量数据、架构原理与实战调优指南

后端消息队列任务调度 【免费下载链接】bullmq BullMQ - Message Queue and Batch processing for NodeJS, Python, .NET, Elixir, Rust and PHP based on Redis or PostgreSQL 项目地址: https://gitcode.com/gh_mirrors/bu/bullmq 点击查看 免费下载 本指南以 Bu…

阅读更多 →
网心云OES Plus刷Armbian教程:拆机短接与系统优化全流程 2026/9/25 3:23:40

网心云OES Plus刷Armbian教程:拆机短接与系统优化全流程

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

阅读更多 →
XAgent ToolServer 深度解析:Manager/Node 双容器架构、完整 API 说明与部署配置实战 2026/9/25 3:23:34

XAgent ToolServer 深度解析:Manager/Node 双容器架构、完整 API 说明与部署配置实战

AI Agent大模型后端任务调度 【免费下载链接】XAgent An Autonomous LLM Agent for Complex Task Solving 项目地址: https://gitcode.com/gh_mirrors/xa/XAgent 点击查看 免费下载 ToolServer 是 XAgent 的工具执行后端:它以 Docker 容器为隔离单元&am…

阅读更多 →
基于 TEN Framework 的 Go 应用性能剖析实战:pprof_app_go 内存剖析方案深度解析 2026/9/25 3:23:34

基于 TEN Framework 的 Go 应用性能剖析实战:pprof_app_go 内存剖析方案深度解析

人工智能AI Agent多模态语音AI 应用 【免费下载链接】ten-framework Open-source framework for conversational voice AI agents 项目地址: https://gitcode.com/TEN-framework/ten-framework 点击查看 免费下载 导读 本文以 TEN Framework 开源仓库中的 Go 示例…

阅读更多 →

今日资讯

本周资讯

本月资讯

看完文章仍有疑问?

联系尧图顾问,获取一对一建站咨询

立即免费咨询 📞 400-888-8888
📞 ✉