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

如何用PubSubEmulatorContainer测试GCP云函数Pub/Sub触发逻辑?

可以用PubSubEmulatorContainer对GCP云函数做集成测试

你的云函数本质是Consumer<PubSubMessage>的实现,测试时无需依赖GCP云函数的部署触发逻辑,直接将这个消费者与Pub/Sub模拟器绑定,模拟消息触发即可。下面是具体实现步骤,同时解决你遇到的NOT_FOUND报错:

一、先解决测试中的NOT_FOUND报错

你遇到的报错核心原因是未在PubSub模拟器中提前创建对应的主题和订阅。模拟器是独立的临时环境,不会自动同步生产环境的资源,必须在测试初始化阶段手动创建。

二、完整集成测试实现

1. 核心依赖配置(以Spring Boot为例)

确保引入Testcontainers PubSub组件和GCP Pub/Sub客户端依赖:

<dependency>
    <groupId>org.testcontainers</groupId>
    <artifactId>gcloud</artifactId>
    <scope>test</scope>
</dependency>
<dependency>
    <groupId>com.google.cloud</groupId>
    <artifactId>spring-cloud-gcp-pubsub</artifactId>
</dependency>

2. 测试类代码实现

import com.google.cloud.pubsub.v1.AckReplyConsumer;
import com.google.cloud.pubsub.v1.Subscriber;
import com.google.pubsub.v1.ProjectSubscriptionName;
import com.google.pubsub.v1.ProjectTopicName;
import com.google.pubsub.v1.PubsubMessage;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.testcontainers.containers.PubSubEmulatorContainer;
import org.testcontainers.junit.jupiter.Container;
import org.testcontainers.junit.jupiter.Testcontainers;
import java.util.concurrent.TimeUnit;
import static org.junit.jupiter.api.Assertions.assertTrue;

@Testcontainers
@SpringBootTest
public class MessageIngestIntegrationTest {

    @Container
    private static final PubSubEmulatorContainer pubSubEmulator = new PubSubEmulatorContainer(
        PubSubEmulatorContainer.DEFAULT_IMAGE_NAME.withTag("latest")
    );

    private static com.google.cloud.pubsub.v1.Publisher publisher;
    private static Subscriber subscriber;

    @Autowired
    private MessageIngest messageIngest;

    @Autowired
    private MessageResource msgResource;

    @Autowired
    private ObjectMapper objectMapper;

    @BeforeAll
    static void setupPubSubResources() throws Exception {
        // 让GCP客户端指向模拟器地址
        System.setProperty("PUBSUB_EMULATOR_HOST", pubSubEmulator.getEmulatorEndpoint());
        String testProjectId = "test-project";
        String testTopicId = "test-topic";
        String testSubscriptionId = "test-subscription";

        // 创建测试主题
        ProjectTopicName topicName = ProjectTopicName.of(testProjectId, testTopicId);
        com.google.cloud.pubsub.v1.TopicAdminClient.create().createTopic(topicName);

        // 创建测试订阅(关联上述主题)
        ProjectSubscriptionName subscriptionName = ProjectSubscriptionName.of(testProjectId, testSubscriptionId);
        com.google.cloud.pubsub.v1.SubscriptionAdminClient.create()
            .createSubscription(subscriptionName, topicName, null, 60);

        // 初始化消息发布者
        publisher = com.google.cloud.pubsub.v1.Publisher.newBuilder(topicName).build();
    }

    @BeforeAll
    void setupSubscriber() {
        // 初始化订阅者,绑定云函数的MessageIngest作为消息处理器
        ProjectSubscriptionName subscriptionName = ProjectSubscriptionName.of("test-project", "test-subscription");
        subscriber = Subscriber.newBuilder(subscriptionName, (PubsubMessage message, AckReplyConsumer consumer) -> {
            try {
                messageIngest.accept(message);
                consumer.ack();
            } catch (Exception e) {
                consumer.nack();
                throw new RuntimeException(e);
            }
        }).build();
        subscriber.startAsync().awaitRunning();
    }

    @AfterAll
    static void cleanUpPubSub() {
        if (publisher != null) {
            publisher.shutdown();
        }
        if (subscriber != null) {
            subscriber.stopAsync().awaitTerminated();
        }
    }

    @Test
    void testMessageTriggerAndProcessing() throws Exception {
        // 构造测试消息
        Message testMsg = new Message("test-001", "sample-content");
        byte[] msgBytes = objectMapper.writeValueAsBytes(testMsg);
        PubsubMessage pubsubMsg = PubsubMessage.newBuilder()
            .setData(com.google.protobuf.ByteString.copyFrom(msgBytes))
            .build();

        // 发布消息到模拟器主题
        publisher.publish(pubsubMsg).get(10, TimeUnit.SECONDS);

        // 等待消息处理完成,验证结果(可替换为CountDownLatch实现更精准的等待)
        TimeUnit.SECONDS.sleep(2);
        assertTrue(msgResource.hasIngested(testMsg));
    }
}

三、关键说明

  • 云函数的MessageIngest是普通Spring Bean,测试时无需依赖GCP云函数的部署机制,直接将它作为订阅者的消息处理器即可模拟触发逻辑。
  • 必须在测试初始化阶段创建主题和订阅,这是解决你NOT_FOUND报错的核心。
  • 若使用Spring Cloud GCP的PubSubSubscriberTemplate,也可替换手动创建的Subscriber,通过模板绑定订阅与MessageIngest消费者。

四、关于contextLoads()的疑问

无需在contextLoads()中实现触发逻辑,而是在@BeforeAll中完成Pub/Sub资源初始化和消费者绑定,再在测试方法中发布消息并验证处理结果。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 16:40:32