ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

PyFlink Table API 连接器使用指南:从 DDL 建表到 Kafka 读写实战

PyFlink Table API 连接器使用指南:从 DDL 建表到 Kafka 读写实战 大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载本篇技术指南围绕 PyFlink Table API 中使用连接器Connector与格式Format的核心议题展开覆盖连接器 jar 包依赖的声明方式、基于execute_sql()的 DDL 建表实践、内置 Source/SinkPandas 互转与from_elements()以及自定义 Source/Sink 的实现路径。读完本文你将掌握在 Python 作业中正确接入 Kafka、JSON 等外部系统的完整方法并能理解连接器在 Python 与 Java/Scala 桥接层上的运行机制。为什么 PyFlink 中的连接器是 jar 包Flink 是一个基于 Java/Scala 构建的分布式计算引擎所有连接器如 Kafka、JDBC、HBase与格式如 JSON、CSV、Avro的实际实现都以 jar 包形式存在。PyFlink 通过 Py4J 网关 与 JVM 侧的TableEnvironment交互因此 Python 侧本身并不重复实现连接器逻辑而是将 jar 包作为作业依赖提交给运行时加载。这意味着在 PyFlink 作业中使用任何第三方连接器前必须先完成两件事准备好对应连接器与格式的 jar 包通过配置或命令行参数将 jar 包声明为作业依赖。连接器的完整列表与通用配置例如支持的 Source/Sink 类型、Kafka 0.10、JDBC、HBase 等请参阅 Table SQL 连接器概览其中明确标注了各连接器支持的读取类型如 Bounded Scan、Unbounded Scan、Lookup与写入类型Streaming Sink、Batch Sink。下载连接器与格式 jar 包通过pipeline.jars指定依赖在 Python 代码中最直接的方式是通过TableConfig的pipeline.jars配置项声明 jar 包列表多个 jar 之间使用分号;分隔table_env.get_config().set(pipeline.jars, file:///my/jar/path/connector.jar;file:///my/jar/path/json.jar)结合依赖管理文档可以得知该配置的详细约束仅支持本地file://协议开头的路径jar 会被上传到集群Windows 环境下路径写法形如file:///E:/my/jar/path/connector.jar若需要将 URL 直接加入客户端与集群两端的 classpath而非上传可改用pipeline.classpaths此时路径必须带协议如file://且需保证 URL 在客户端与集群上均可访问table_env.get_config().set(pipeline.classpaths, file:///my/jar/path/connector.jar;file:///my/jar/path/udf.jar)通过命令行参数提交除了在代码中配置也可以在提交作业时使用--jarfile命令行参数指定 jar。需要注意该参数只支持指定一个jar 文件若存在多个 jar 依赖需要预先打成 fat jar 再提交。Python DataStream API 中的等价方式如果你在同一个作业中混用 Python DataStream API 与 Python Table API依赖应通过 DataStream API 统一声明以保证两端同时生效。DataStream API 提供如下方法stream_execution_environment.add_jars(file:///my/jar/path/connector1.jar, file:///my/jar/path/connector2.jar) stream_execution_environment.add_classpaths(file:///my/jar/path/connector1.jar, file:///my/jar/path/connector2.jar)如何通过 DDL 使用连接器在 PyFlink Table API 中推荐使用 DDL 定义 source 和 sink。DDL 通过TableEnvironment.execute_sql()方法执行——从 table_environment.py 的源码可以看到execute_sql()支持 DDL/DML/DQL/SHOW/DESCRIBE/EXPLAIN/USE 等语句对 DDL/DCL 语句操作完成后即返回结果对 DML/DQL 语句作业提交后返回TableResult。定义好表之后就可以在 SQL 查询或 Table API 中直接引用该表了。下面是一个同时定义 Kafka source 与 sink 的 DDL 示例source_ddl CREATE TABLE source_table( a VARCHAR, b INT ) WITH ( connector kafka, topic source_topic, properties.bootstrap.servers kafka:9092, properties.group.id test_3, scan.startup.mode latest-offset, format json ) sink_ddl CREATE TABLE sink_table( a VARCHAR ) WITH ( connector kafka, topic sink_topic, properties.bootstrap.servers kafka:9092, format json ) t_env.execute_sql(source_ddl) t_env.execute_sql(sink_ddl) t_env.sql_query(SELECT a FROM source_table) \ .execute_insert(sink_table).wait()其中各 WITH 参数的含义如下参数作用connector连接器工厂标识Flink 通过该标识如kafka、jdbc在 SPI 工厂中查找对应的 TableSource/TableSink 实现topicKafka topic 名称source 与 sink 各自独立指定properties.bootstrap.serversKafka broker 地址列表properties.group.id消费者组 ID仅 source 需要scan.startup.mode消费起始位置如earliest-offset、latest-offset常见流式消费场景使用latest-offsetformat消息体格式如json、csv需搭配对应格式 jar关于connector与format如何被解析Table SQL 连接器概览 中给出了底层机制连接属性被转换为字符串键值对Flink 通过 Java 的 Service Provider InterfaceSPI机制加载所有可用的工厂Factory再依据工厂标识精确匹配唯一一个工厂来创建 source、sink 与格式。若找不到工厂或多个工厂同时匹配会抛出异常并附带被考虑过的工厂与支持的属性信息。提示官方示例原文中的connector kafka属于笔误正确写法为connector kafka上例已修正。DDL 中任何键值格式错误都会在execute_sql()阶段直接报错。完整示例PyFlink 中 Kafka source/sink 与 JSON 格式下面是在 PyFlink 中使用 Kafka source/sink 与 JSON 格式的完整可运行示例from pyflink.table import TableEnvironment, EnvironmentSettings def log_processing(): env_settings EnvironmentSettings.in_streaming_mode() t_env TableEnvironment.create(env_settings) # specify connector and format jars t_env.get_config().set(pipeline.jars, file:///my/jar/path/connector.jar;file:///my/jar/path/json.jar) source_ddl CREATE TABLE source_table( a VARCHAR, b INT ) WITH ( connector kafka, topic source_topic, properties.bootstrap.servers kafka:9092, properties.group.id test_3, scan.startup.mode latest-offset, format json ) sink_ddl CREATE TABLE sink_table( a VARCHAR ) WITH ( connector kafka, topic sink_topic, properties.bootstrap.servers kafka:9092, format json ) t_env.execute_sql(source_ddl) t_env.execute_sql(sink_ddl) t_env.sql_query(SELECT a FROM source_table) \ .execute_insert(sink_table).wait() if __name__ __main__: log_processing()示例中的关键调用链在源码中均有对应实现EnvironmentSettings.in_streaming_mode()创建流处理模式的环境设置该设置仅在TableEnvironment实例化时生效创建后不可变更execute_sql()内部委托给 JVM 侧的executeSql(stmt)返回TableResult见 table_environment.pysql_query()通过_j_tenv.sqlQuery(query)求值 SQL 并包装为 Python 侧Table见 table_environment.pyexecute_insert()将表写入已注册的 sink并返回TableResult见 table.py调用.wait()会阻塞直到作业执行完成。内置的 Sources 和 Sinks有些 source 和 sink 被内置在 Flink 中无需额外 jar 即可直接使用。内置 source 包括将 Pandas DataFrame 作为数据源、将元素集合作为数据源内置 sink 包括将表数据转换为 Pandas DataFrame 等。这些能力由 Python 侧原生实现并通过 JVM 网关桥接是 PyFlink 区别于纯 SQL 连接器场景的便捷特性。与 Pandas 之间互转PyFlink 表支持与 Pandas DataFrame 之间互相转换。从from_pandas的源码table_environment.py可以看到该方法底层依赖 PyArrow 将 DataFrame 转换为 Arrow schema再构造 RowTypesplits_num参数决定数据被切分的份数进而影响并行 source 任务的数目不指定时使用默认并行度。而to_pandastable.py通过ArrowUtils.collectAsPandasDataFrame在客户端收集结果并转换为 DataFrame因此调用前请确保表内容能够放进客户端内存。from pyflink.table.expressions import col import pandas as pd import numpy as np # 创建一个 PyFlink 表 pdf pd.DataFrame(np.random.rand(1000, 2)) table t_env.from_pandas(pdf, [a, b]).filter(col(a) 0.5) # 将 PyFlink 表转换成 Pandas DataFrame pdf table.to_pandas()from_pandas的第二个参数schema支持三种形态字段名字符串列表如[a, b]、字段类型列表如[DataTypes.DOUBLE(), DataTypes.DOUBLE()]、或完整的DataTypes.ROW(...)行类型均用于指定转换后的表结构。from_elements()from_elements()用于从一个元素集合中创建一张表。元素类型必须是可支持的原子类型或复杂类型。根据源码注释内置可支持的原子类型包括int、long、str、unicode、bool、float、bytearray、datetime.date、datetime.time、datetime.datetime、datetime.timedelta、decimal.Decimal可支持的复合类型包括list、tuple、dict、array与pyflink.table.Row。from pyflink.table import DataTypes table_env.from_elements([(1, Hi), (2, Hello)]) # 使用第二个参数指定自定义字段名 table_env.from_elements([(1, Hi), (2, Hello)], [a, b]) # 使用第二个参数指定自定义表结构 table_env.from_elements([(1, Hi), (2, Hello)], DataTypes.ROW([DataTypes.FIELD(a, DataTypes.INT()), DataTypes.FIELD(b, DataTypes.STRING())]))不指定 schema 时复合类型会被自动解包字段名自动生成为_1、_2。以上查询返回的表如下----------- | a | b | | 1 | Hi | ----------- | 2 | Hello | -----------from_elements()还支持第三个参数verify_schema默认为True用于控制是否对元素做 schema 校验同时支持传入Expression元素如row(1, abc, 2.0)直接构建表见 table_environment.py。从实现看元素最终会被序列化到临时文件再由 JVM 侧的PythonTableUtils.createTableFromElement读取构建表对象。仓库中的单元测试如 test_calc.py大量使用from_elements配合 DDL 注册的测试 sink 来验证查询逻辑可作为该 API 实际用法的参考。此外PyFlink 还保留了基于TableSource的编程式 API例如 sources.py 中的CsvTableSource可通过 builder 链式配置字段名、字段类型、分隔符、引号字符、是否跳过首行、是否宽松解析等参数作为 DDL 之外创建 source 的补充手段。用户自定义的 source 和 sink在某些情况下你可能想要自定义 source 或 sink。目前source 和 sink 必须使用 Java/Scala 实现你可以定义一个TableFactory然后通过 DDL 在 PyFlink 作业中引用它。其核心原理与内置连接器一致——DDL 中声明的connector标识会通过 SPI 工厂查找机制定位到你的自定义实现。更完整的实现指南包括DynamicTableSourceFactory/DynamicTableSinkFactory的接口约定与注册流程可参阅 用户自定义 Sources Sinks 的 Java/Scala 文档。总结与注意事项在 PyFlink 中接入外部系统核心路径可以概括为三步声明 jar 依赖 → 用 DDL 定义 source/sink → 在 SQL 或 Table API 中读写数据。同时需要留意以下几点连接器与格式必须作为 jar 依赖显式声明pipeline.jars仅支持file://本地路径多个 jar 用;分隔DDL 的WITH参数全部为字符串键值对connector与format是工厂匹配的关键标识写错会导致工厂查找失败内置的 Pandas 互转与from_elements()不需要额外 jar适合本地调试与快速原型验证to_pandas()会一次性将结果收集到客户端只适用于数据量可放进内存的场景自定义 source/sink 需要以 Java/Scala 实现并通过 SPI 工厂注册PyFlink 侧通过 DDL 透明使用。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐PyFlink Table API 连接器使用指南DDL 定义、依赖管理、预置 Source/Sink 与 Pandas 互转PyFlink Table API 连接器使用指南DDL 定义、依赖管理、预置 Source/Sink 与 Pandas 互转 本文基于 Apache Fli大数据流处理批处理数据工程Delta Kernel 连接器开发实战用纯 Java API 读写 Delta Lake 表的完整指南Delta Kernel 连接器开发实战用纯 Java API 读写 Delta Lake 表的完整指南 导读 Delta Kernel 是 Delta La湖仓一体数据工程数据湖Flink Table API 实时报表实战从 Kafka 到 MySQL 再到 Grafana 的端到端看板搭建Flink Table API 实时报表实战从 Kafka 到 MySQL 再到 Grafana 的端到端看板搭建 Apache Flink 的 Table大数据流处理批处理数据工程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表