Spark 2.3.1单源多Sink流处理异常:仅一个Sink正常工作
嘿,我之前在Spark 2.3.x版本里也踩过一模一样的多Sink坑!先给你说清楚为啥会出现你说的“一个流正常、另一个没输出,注释掉正常的那个另一个才生效”的情况:
在Spark 2.3.1里,每调用一次writeStream.start(),就会创建一个独立的流查询。这些查询会各自从数据源拉取数据——但像你示例里的socket源,它本质是单消费者的,第一个启动的查询会独占socket连接,第二个查询根本拿不到任何数据。哪怕换成Kafka这类多消费者源,也可能因为Checkpoint的状态管理、偏移量同步的问题,导致只有一个查询能正常消费,另一个一直蹲在那等数据。
更关键的是,Spark 2.3.x确实没有后续版本才有的foreachBatch这个多Sink神器,所以得换个思路绕过去。
几个可行的解决方案
1. 写个自定义ForeachWriter,一次处理两个Sink
这是最直接的办法:自己实现一个ForeachWriter,在同一个批次里同时把数据写入HDFS和Kafka。这样整个流程只有一个流查询,自然不会出现抢数据源的问题。
给你个大概的代码示例参考:
import org.apache.spark.sql.ForeachWriter import org.apache.spark.sql.Row import org.apache.kafka.clients.producer.{KafkaProducer, ProducerRecord} import org.apache.hadoop.fs.Path class MultiSinkWriter extends ForeachWriter[Row] { private var kafkaProducer: KafkaProducer[String, String] = _ private var tempParquetPath: String = _ override def open(partitionId: Long, version: Long): Boolean = { // 初始化Kafka生产者 val kafkaProps = new java.util.Properties() kafkaProps.put("bootstrap.servers", "host1:port1,host2:port2") kafkaProps.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer") kafkaProps.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer") kafkaProducer = new KafkaProducer[String, String](kafkaProps) // 生成临时路径存当前批次的Parquet,避免写入过程中数据被读取 tempParquetPath = s"path/to/destination/dir/temp_batch_${partitionId}_${version}" true } override def process(row: Row): Unit = { val data = row.getString(0) // 写Kafka kafkaProducer.send(new ProducerRecord[String, String]("updates", data)) // 这里可以用Spark的API把当前批次的数据写入临时Parquet路径,或者直接用ParquetWriter API // 注意:如果用DataFrame的话,要保证是当前批次的小批量数据,避免内存溢出 } override def close(errorOrNull: Throwable): Unit = { // 关闭Kafka生产者前先刷数据 if (kafkaProducer != null) { kafkaProducer.flush() kafkaProducer.close() } // 把临时Parquet文件原子移动到正式目录,保证数据一致性 val fs = org.apache.hadoop.fs.FileSystem.get(new java.net.URI("hdfs://your-cluster"), new org.apache.hadoop.conf.Configuration()) val tempPath = new Path(tempParquetPath) val destPath = new Path("path/to/destination/dir") if (fs.exists(tempPath)) { fs.rename(tempPath, destPath) } } } // 用这个自定义Writer替代两个独立的writeStream words.writeStream .foreach(new MultiSinkWriter()) .option("checkpointLocation", "shared_checkpoint_path") .start() spark.streams.awaitAnyTermination()
提醒一下:这个方案需要你自己处理事务性哦,如果写入Kafka成功但HDFS失败,得考虑怎么回滚,避免数据不一致。如果对一致性要求不高,这个方案足够简单好用。
2. 先写中间存储,再分流转储
如果自定义Writer嫌麻烦,你可以先把流数据写入一个中间的多消费者存储(比如Kafka),然后启动两个独立的流:一个从中间Kafka读了写HDFS,另一个读了写目标Kafka。
流程大概是:
- 源数据 → 中间Kafka主题
stream_temp - 流1:从
stream_temp读 → 写入HDFS Parquet - 流2:从
stream_temp读 → 写入目标Kafka主题updates
这个方案的好处是不用写自定义代码,每个流都是独立的,出错了可以单独重启。缺点是多了一次写入,会有额外的延迟和存储开销。
3. 缓存流数据(仅适用于多消费者数据源)
如果你一定要用两个独立的writeStream,可以先把流的Dataset缓存起来,避免重复计算:
val words: Dataset[String] = lines.as[String] .flatMap(_.split(" ")) .map(_.concat(Random.nextInt(10).toString)) .cache() // 缓存计算结果,避免两个流重复处理 // 写入HDFS的流 words.writeStream .format("parquet") .option("path", "path/to/destination/dir") .option("checkpointLocation", "some_location_1") .start() // 写入Kafka的流 words.writeStream .format("kafka") .option("kafka.bootstrap.servers", "host1:port1,host2:port2") .option("topic", "updates") .option("checkpointLocation", "some_location_2") .start() spark.streams.awaitAnyTermination()
但注意哦:这个方法只能解决重复计算的问题,像socket这种单消费者源,还是会出现只有一个流能拿到数据的情况。所以只适用于Kafka这类支持多消费者的数据源,而且要保证两个流的Checkpoint路径完全独立。
针对你简化测试的问题(Console正常,自定义ForeachWriter无输出)
这种情况大概率是自定义Writer里藏了小问题:要么open方法返回了false,要么process方法里有未捕获的异常。你可以在Writer的各个方法里加打印日志,比如:
- 在
open里加println(s"打开分区$partitionId,版本$version") - 在
process里加println(s"处理数据:$row") - 在
close里加println(s"关闭时的错误:${errorOrNull.getMessage}")
这样就能快速定位是哪一步出问题了。
内容的提问来源于stack exchange,提问作者bitan

