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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.01 03:16:04