Scala+Spark Structured Streaming向WSL Kafka推送数据报awaitResult异常
问题原因
- 流任务启动后未阻塞主线程
你在调用.start()启动Structured Streaming任务后没有调用.awaitTermination(),主线程执行完所有代码直接退出,JVM销毁Spark上下文,就会触发RpcEnvStoppedException异常。 - Kafka写入格式不符合要求
Spark Structured Streaming写入Kafka时,要求DataFrame必须包含value列(可选包含key列),且列类型为字符串或二进制类型。你当前直接写入读取到的CSV全字段,不符合Kafka写入要求,所以消息无法正常发送。 - WSL环境Kafka网络配置错误
默认部署在WSL内的Kafka只监听WSL内部的本地地址,Windows主机通过localhost:9092无法正常连接,就算代码无错也无法发送消息到WSL的Kafka集群。 - 定义的Schema和实际CSV数据不匹配
你在Schema中将price字段定义为Integer类型,但实际生成的CSV数据中price是$100,000、Priceless这类字符串,类型转换失败会导致流任务异常终止。 - CSV流式读取配置问题
Structured Streaming默认只读取流式任务启动后新增到目录的文件,你是先生成CSV文件再启动流任务,默认不会读取已存在的文件。
解决方案
第一部分:代码修改
- 修正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() // 新增这行,阻塞主线程直到流任务终止
- 修正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")
- 修正CSV读取配置,允许读取已有文件
val csvDF = spark .readStream .option("sep", ",") .option("startingOffsets", "earliest") // 新增这行,读取目录下已有的历史文件 .schema(userSchema) .format("csv") .load("input/") // 统一用/作为路径分隔符,Windows也能正常识别
- 修正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服务即可。
第三部分:验证步骤
- 先在WSL中启动Kafka消费者,监听目标topic:
kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic target_topic --from-beginning
- 再运行IDEA中的Scala代码,即可看到消费者接收到JSON格式的消息。
内容的提问来源于stack exchange,提问作者Angel Torres
相关产品推荐
相关产品推荐

