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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 04:18:01