You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PyFlink 1.16中如何将Protobuf描述自动转换为Flink Schema?

PyFlink 1.16对Protobuf的处理方式相比1.10有较大调整,原ProtobufSchemaConverter的Python实现已被移除。以下两种方案可解决自动生成Flink Schema的需求:

方案一:使用Table API的Protobuf Format(推荐)

这是官方推荐的方式,无需手动编写Schema,直接通过配置protobuf格式参数让Flink自动推断Schema。

示例1:基于生成的Protobuf类

假设你已将.proto文件编译为Java类(PyFlink的Protobuf格式依赖Java类或描述符文件),可通过类全名指定消息类型:

from pyflink.table import TableEnvironment, EnvironmentSettings

env_settings = EnvironmentSettings.in_streaming_mode()
t_env = TableEnvironment.create(env_settings)

# 自动推断Schema的表DDL
t_env.execute_sql("""
    CREATE TABLE proto_kafka_source (
        -- 无需手动定义字段,Flink会从Protobuf类自动推断
    ) WITH (
        'connector' = 'kafka',
        'topic' = 'test_topic',
        'properties.bootstrap.servers' = 'localhost:9092',
        'format' = 'protobuf',
        'protobuf.message-class-name' = 'com.yourpackage.test' -- 替换为你的Protobuf类全限定名
    )
""")

# 验证Schema
t_env.from_path("proto_kafka_source").print_schema()

示例2:基于Protobuf描述符文件

若不想编译为Java类,可使用.desc格式的描述符文件:

t_env.execute_sql("""
    CREATE TABLE proto_kafka_source (
        -- 无需手动定义字段
    ) WITH (
        'connector' = 'kafka',
        'topic' = 'test_topic',
        'properties.bootstrap.servers' = 'localhost:9092',
        'format' = 'protobuf',
        'protobuf.descriptor-file' = '/path/to/your/proto_desc.desc', -- 描述符文件路径
        'protobuf.message-name' = 'test' -- 指定消息名称
    )
""")

方案二:程序化生成Schema对象

若你需要在Python代码中显式获取Schema对象(如动态创建表),可通过PyFlink的Java网关调用Java版的ProtobufSchemaConverter:

步骤1:获取Protobuf描述符

首先从.desc文件加载描述符:

from google.protobuf.descriptor_pb2 import FileDescriptorSet

# 读取描述符文件
with open('/path/to/your/proto_desc.desc', 'rb') as f:
    fds = FileDescriptorSet()
    fds.ParseFromString(f.read())

# 定位到test消息的描述符
test_msg_desc = None
for file_desc in fds.file:
    for msg in file_desc.message_type:
        if msg.name == 'test':
            test_msg_desc = msg
            break
    if test_msg_desc:
        break
from pyflink.java_gateway import get_gateway
from pyflink.table import Schema

# 获取Java网关实例
gateway = get_gateway()

# 调用Java版ProtobufSchemaConverter
JavaProtobufConverter = gateway.jvm.org.apache.flink.table.protobuf.ProtobufSchemaConverter
java_schema = JavaProtobufConverter.fromDescriptor(test_msg_desc)

# 转换为Python Schema对象
python_schema = Schema(java_schema)

# 验证Schema结构
print(python_schema)

步骤3:使用生成的Schema创建表

t_env.create_table(
    "dynamic_proto_source",
    python_schema,
    connector_properties={
        'connector': 'kafka',
        'topic': 'test_topic',
        'properties.bootstrap.servers': 'localhost:9092',
        'format': 'protobuf',
        'protobuf.descriptor-file': '/path/to/your/proto_desc.desc',
        'protobuf.message-name': 'test'
    }
)

注意事项

  • 确保PyFlink环境包含Protobuf依赖:可通过pip install apache-flink[protobuf]安装,或提交作业时指定flink-protobuf Jar包。
  • 若使用Protobuf类,需确保类文件在作业的类路径中(如打包进Fat Jar或通过--jar参数指定)。

内容的提问来源于stack exchange,提问作者Lost_Deviation

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.02 09:05:14