PySpark无法订阅Kafka Topic问题:程序直接退出无流数据输出
核心问题原因
你使用Spark Structured Streaming对接Kafka时仅完成了流式DataFrame的Schema定义,没有启动实际的流计算任务,也未设置进程阻塞机制,代码执行完打印语句后主进程直接退出,因此无法持续消费Kafka消息。
解决方案
1. 修正业务代码
补充流查询启动逻辑和进程阻塞逻辑,完整修正后代码如下:
import findspark findspark.init("/usr/local/spark-3.1.2-bin-hadoop2.7") from pyspark.sql import SparkSession KAFKA_TOPIC = "kafka-spark" KAFKA_SERVER = "localhost:9092" # 创建SparkSession实例 spark_session = SparkSession.builder.appName("KafkaSparkDemo").getOrCreate() # 订阅Kafka Topic df = spark_session \ .readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", KAFKA_SERVER) \ .option("subscribe", KAFKA_TOPIC) \ # 若需要消费Topic历史消息可开启以下配置,默认仅消费启动后新产生的消息 # .option("startingOffsets", "earliest") \ .load() # 转换key、value字段为字符串格式 result_df = df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") # 定义控制台输出Sink并启动流查询 query = result_df.writeStream \ .format("console") \ .outputMode("append") \ .start() # 阻塞主进程,持续运行流任务直到手动终止 query.awaitTermination()
2. 优化提交命令
你当前使用Structured Streaming API,无需引入Spark Streaming的冗余依赖包,优化后提交命令如下:
spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.1.2 /usr/local/spark-3.1.2-bin-hadoop2.7/examples/src/main/python/kafkaspark.py
补充说明
Spark Structured Streaming为懒执行模式,仅定义流式DataFrame不会触发实际计算,必须调用start()方法启动流任务,同时通过awaitTermination()阻止主进程退出,才能持续消费Kafka中的数据并输出到控制台。
内容的提问来源于stack exchange,提问作者haq
相关产品推荐
相关产品推荐

