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

使用Flink Table API读取Kafka Topic时出错,求解决

解决Flink读取Kafka流并打印数据的问题

错误核心原因

从报错栈的java.lang.NoSuchMethodError可以看出,Flink Kafka连接器与Flink核心依赖版本不兼容,同时代码中print()的调用方式也存在问题。

可行解决方法

1. 确保依赖版本完全匹配

  • 必须使用和Flink版本(1.17.0)完全一致的flink-sql-connector-kafka jar包,从官方渠道下载对应版本,避免使用第三方或版本不匹配的包。
  • 清理本地环境中其他版本的Flink相关依赖,防止冲突。

2. 修正Table API的打印调用方式

TableResult.print()方法本身会直接输出查询结果,不需要再用Python的print()包裹:

# 替换原代码中的print(table_result.print())
table_result = tbl_env.execute_sql("SELECT * FROM sales_usd")
table_result.print()

3. 改用DataStream API打印数据(更直观的调试方式)

将Table转换为DataStream后调用print(),适合流式场景的调试:

# 在定义tbl之后添加以下代码
from pyflink.datastream import DataStream

# 将Table转换为DataStream
ds: DataStream = tbl_env.to_data_stream(tbl)
# 打印流数据
ds.print()
# 触发任务执行
env.execute("Print Kafka Sales Data")

4. 检查集群与本地依赖一致性

如果是连接远程Flink集群运行任务,确保集群中已部署相同版本的Kafka连接器,或者通过pipeline.jars参数正确上传依赖jar包,避免集群端依赖缺失。

完整修正后的代码示例(Table API方式)

import os

from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment, EnvironmentSettings
from pyflink.table.expressions import *


def main():
    # Create streaming environment
    env = StreamExecutionEnvironment.get_execution_environment()
    settings = EnvironmentSettings.in_streaming_mode()

    # create table environment
    tbl_env = StreamTableEnvironment.create(stream_execution_environment=env,
                                            environment_settings=settings)

    # add kafka connector dependency
    kafka_jar = os.path.join(os.path.abspath(os.path.dirname(__file__)),
                            'flink-sql-connector-kafka-1.17.0.jar')

    tbl_env.get_config()\
            .get_configuration()\
            .set_string("pipeline.jars", "file://{}".format(kafka_jar))

    #######################################################################
    # Create Kafka Source Table with DDL
    #######################################################################
    src_ddl = """
        CREATE TABLE sales_usd (
            seller_id VARCHAR,
            amount_usd DOUBLE,
            sale_ts BIGINT
        ) WITH (
            'connector' = 'kafka',
            'topic' = 'sales-usd',
            'properties.bootstrap.servers' = 'localhost:9092',
            'properties.group.id' = 'sales-1',
            'scan.startup.mode' = 'earliest-offset',
            'format' = 'json'
        )
    """

    tbl_env.execute_sql(src_ddl)

    # create and initiate loading of source Table
    tbl = tbl_env.from_path('sales_usd').select(col("seller_id"),col("amount_usd"),col("sale_ts"))

    print('\nSource Schema')
    tbl.print_schema()
    
    # 正确打印流式数据
    table_result = tbl_env.execute_sql("SELECT * FROM sales_usd")
    table_result.print()


if __name__ == "__main__":
    main()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 19:54:57