使用Flink Table API读取Kafka Topic时出错,求解决
解决Flink读取Kafka流并打印数据的问题
错误核心原因
从报错栈的java.lang.NoSuchMethodError可以看出,Flink Kafka连接器与Flink核心依赖版本不兼容,同时代码中print()的调用方式也存在问题。
可行解决方法
1. 确保依赖版本完全匹配
- 必须使用和Flink版本(1.17.0)完全一致的
flink-sql-connector-kafkajar包,从官方渠道下载对应版本,避免使用第三方或版本不匹配的包。 - 清理本地环境中其他版本的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
相关产品推荐
相关产品推荐

