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

Google Pub/Sub订阅者无法接收消息(Akka Alpakka环境)

Troubleshooting Akka Alpakka Google Pub/Sub Subscription Not Receiving Messages

I’ve been in your shoes before—when publishing works but subscribing goes silent, it’s frustrating especially when you’ve already confirmed credentials and permissions are solid. Let’s break down the most likely fixes for your issue:

1. You Forgot to Start the Stream

Akka Streams doesn’t run automatically just by defining the graph. Looking at your code, you’ve wired up the source to the sink with .to(ackSink), but this only creates a RunnableGraph—you need to explicitly call .run() to kick off the stream execution. This explains why GooglePubSubSource.createLogic never runs.

Update your code’s final section like this:

subscriptionSource
 .map { message =>
   val data = message.message.data
   println(s"received a message: $data")
   message.ackId
 }
 .groupedWithin(1000, 1.minute)
 .map(AcknowledgeRequest.apply)
 .to(ackSink)
 .run() // Critical: This starts the stream processing

2. Adjust Polling Settings for Testing

By default, the Pub/Sub source might have a longer poll interval that makes it seem like no messages are coming in. Try overriding the settings to pull more frequently during debugging:

val subscribeSettings = SubscribeSettings.default
  .withPollInterval(1.second) // Shorten interval to see results faster
  .withMaxMessagesPerPoll(10) // Adjust batch size if needed

val subscriptionSource: Source[ReceivedMessage, NotUsed] = 
  GooglePubSub.subscribe(
    projectId, 
    apiKey, 
    clientEmail, 
    privateKey, 
    subscription,
    subscribeSettings
  )

3. Ensure the Program Doesn’t Exit Prematurely

When using the App trait, the JVM might shut down before the stream has a chance to initialize. Add code to block the main thread until the stream completes:

import scala.concurrent.Await
import scala.concurrent.duration._

val streamCompletion: Future[Done] = subscriptionSource
 .map { message =>
   val data = message.message.data
   println(s"received a message: $data")
   message.ackId
 }
 .groupedWithin(1000, 1.minute)
 .map(AcknowledgeRequest.apply)
 .to(ackSink)
 .run()

Await.result(streamCompletion, Duration.Inf) // Keep the app running until the stream ends

4. Add Debug Logs to Catch Hidden Errors

Enable Akka’s debug logging to see what’s happening under the hood. Create an application.conf file in your resources folder with:

akka {
  loglevel = "DEBUG"
  logger = "akka.event.slf4j.Slf4jLogger"
}

This will show you stream initialization logs, Pub/Sub API call details, and any subtle errors that might not be visible in your print statements.

5. Double-Check Subscription-Topic Mapping

Even with full permissions, it’s easy to mix up subscriptions and topics. Confirm:

  • The subscription somesubscription is linked to the exact topic you’re publishing messages to
  • The subscription exists in your weirdproject project (verify via Google Cloud Console)

Start with adding .run() to your stream—it’s the most likely fix for your createLogic not executing issue. Let me know if you still hit snags after trying these steps!

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 07:44:58