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

Scala+Spark Structured Streaming向WSL Kafka推送数据报awaitResult异常

问题原因

  1. 流任务启动后未阻塞主线程
    你在调用.start()启动Structured Streaming任务后没有调用.awaitTermination(),主线程执行完所有代码直接退出,JVM销毁Spark上下文,就会触发RpcEnvStoppedException异常。
  2. Kafka写入格式不符合要求
    Spark Structured Streaming写入Kafka时,要求DataFrame必须包含value列(可选包含key列),且列类型为字符串或二进制类型。你当前直接写入读取到的CSV全字段,不符合Kafka写入要求,所以消息无法正常发送。
  3. WSL环境Kafka网络配置错误
    默认部署在WSL内的Kafka只监听WSL内部的本地地址,Windows主机通过localhost:9092无法正常连接,就算代码无错也无法发送消息到WSL的Kafka集群。
  4. 定义的Schema和实际CSV数据不匹配
    你在Schema中将price字段定义为Integer类型,但实际生成的CSV数据中price是$100,000、Priceless这类字符串,类型转换失败会导致流任务异常终止。
  5. CSV流式读取配置问题
    Structured Streaming默认只读取流式任务启动后新增到目录的文件,你是先生成CSV文件再启动流任务,默认不会读取已存在的文件。

解决方案

第一部分:代码修改

  1. 修正Kafka写入逻辑并添加主线程阻塞
// 把所有字段转成JSON字符串作为Kafka的value,符合写入要求
val query = csvDF
  .select(
    col("order_id").cast("string").as("key"), // 可选指定key,这里用订单id作为消息key
    functions.to_json(functions.struct("*")).cast("string").as("value")
  )
  .writeStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "localhost:9092")
  .option("topic", "target_topic")
  .option("checkpointLocation", "tmp/vaquarkhan/checkpoint")
  .start()

query.awaitTermination() // 新增这行,阻塞主线程直到流任务终止
  1. 修正Schema的price字段类型
val userSchema = new StructType().add("order_id", "integer").add("customer_id", "integer").add("customer_name", "string")
  .add("product_id", "integer").add("product_name", "string").add("product_category", "string").add("qty", "integer")
  .add("price", "string") // 从integer改为string,匹配实际数据格式
  .add("datetime", "string").add("country", "string").add("city", "string")
  .add("ecommerce_website_name", "string").add("payment_txn_id", "string")
  .add("payment_txn_success", "string").add("failure_reason", "string")
  1. 修正CSV读取配置,允许读取已有文件
val csvDF = spark
  .readStream
  .option("sep", ",")
  .option("startingOffsets", "earliest") // 新增这行,读取目录下已有的历史文件
  .schema(userSchema)
  .format("csv")
  .load("input/") // 统一用/作为路径分隔符,Windows也能正常识别
  1. 修正CSV文件生成逻辑
    写完CSV内容后关闭输出流,确保内容写入磁盘:
for(i <- 1 to 2000){
  // 原有生成transaction的逻辑保持不变
}
printWriter.close() // 新增这行,关闭流保证文件写入完成

第二部分:WSL Kafka配置修改

修改WSL中Kafka的config/server.properties文件:

listeners=PLAINTEXT://0.0.0.0:9092 # 监听所有网卡地址,允许外部访问
advertised.listeners=PLAINTEXT://localhost:9092 # 对外暴露的访问地址为localhost,适配Windows主机访问

修改完成后重启Kafka broker服务即可。

第三部分:验证步骤

  1. 先在WSL中启动Kafka消费者,监听目标topic:
kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic target_topic --from-beginning
  1. 再运行IDEA中的Scala代码,即可看到消费者接收到JSON格式的消息。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 05:57:03