Spark使用Pub/Sub Lite发布消息报分区计数刷新失败异常
Spark Structured Streaming 发布消息到GCP Pub/Sub Lite 报分区计数刷新失败
问题背景
业务场景下需要在Spark Structured Streaming的forEachBatch sink中自定义消息发布逻辑,无法直接使用官方内置的writeStream Pub/Sub连接器,因此采用forEachPartition嵌套forEach的实现方案,在forEach算子中逐行处理DataFrame数据完成消息发布。目前部分消息可正常发布,部分场景下抛出如下异常:
2022-06-07 10:08:17 WARN PartitionCountWatcherImpl:101 - Failed to refresh partition count com.google.api.gax.rpc.ApiException: at com.google.cloud.pubsublite.internal.CheckedApiException.<init>(CheckedApiException.java:51) at com.google.cloud.pubsublite.internal.CheckedApiException.<init>(CheckedApiException.java:55) at com.google.cloud.pubsublite.internal.ExtractStatus.toCanonical(ExtractStatus.java:49) at com.google.cloud.pubsublite.internal.wire.PartitionCountWatcherImpl.pollTopicConfig(PartitionCountWatcherImpl.java:92) at com.google.cloud.pubsublite.internal.wire.PartitionCountWatcherImpl.onAlarm(PartitionCountWatcherImpl.java:71) at com.google.cloud.pubsublite.internal.AlarmFactory.lambda$null$0(AlarmFactory.java:41) at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) at java.util.concurrent.FutureTask.runAndReset(FutureTask.java:308) at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$301(ScheduledThreadPoolExecutor.java:180) at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:294) at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) at java.lang.Thread.run(Thread.java:748) Caused by: java.lang.InterruptedException at com.google.common.util.concurrent.AbstractFuture.get(AbstractFuture.java:456) at com.google.common.util.concurrent.FluentFuture$TrustedFuture.get(FluentFuture.java:100) at com.google.common.util.concurrent.ForwardingFuture.get(ForwardingFuture.java:73) at com.google.cloud.pubsublite.internal.wire.PartitionCountWatcherImpl.pollTopicConfig(PartitionCountWatcherImpl.java:81) ... 9 more
根因说明
- 该异常的核心触发原因是日志最底层的
java.lang.InterruptedException:Pub/Sub Lite发布客户端内置了后台定时调度线程PartitionCountWatcherImpl,会周期性拉取Topic的最新分区数用于消息路由,当这个后台线程被强制中断时就会抛出该WARN级别的日志。 - 最常见的触发场景是客户端实例生命周期管理错误:如果在逐行执行的
forEach算子内部创建Pub/Sub Lite发布客户端,每处理一行数据就新建一个客户端实例,会生成大量冗余后台调度线程,这些临时客户端在单行数据处理完成后被JVM回收、或者Spark任务线程被销毁时,绑定的后台线程就会被中断抛出异常。 - 偶发单条该WARN属于正常现象:
PartitionCountWatcherImpl自带失败重试逻辑,单次轮询被中断不会影响整体功能,下次定时轮询会自动恢复。如果该日志高频出现且伴随消息发布失败、发布延迟飙升,才需要针对性修复。
修复方案
- 调整客户端初始化位置:Pub/Sub Lite发布客户端是重资源长连接对象,必须在
forEachPartition的分区入口处初始化,保证单个Spark分区内的所有行数据复用同一个客户端实例,禁止在逐行处理的forEach逻辑内创建客户端。 - 做好资源优雅回收:如果是分区级独占的客户端实例,在分区内所有数据处理完成后,必须在
finally块中显式调用客户端的关闭方法,等待后台线程正常退出后再释放资源;如果是Executor级静态复用的客户端,不需要随分区关闭,避免重复创建销毁的开销。
参考实现逻辑(Scala):val query = df.writeStream.foreachBatch { (batchDf, batchId) => batchDf.foreachPartition { rowIterator => // 分区维度初始化一次客户端 val publisher = initializePubSubLitePublisher() try { while (rowIterator.hasNext) { val row = rowIterator.next() val pubsubMessage = convertRowToMessage(row) // 复用客户端发布消息 publisher.publish(pubsubMessage) } } finally { // 分区处理完成后优雅关闭客户端,中断后台线程 publisher.shutdown() publisher.awaitTermination(30, TimeUnit.SECONDS) } } }.start() - 若日志偶发且消息发布成功率符合预期,不需要额外处理,可通过日志配置调整该类WARN日志的打印级别,减少无效日志干扰。
内容的提问来源于stack exchange,提问作者Pranjal Singh
相关产品推荐
相关产品推荐

