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

编写PubSub消息发布方法的JUnit测试用例时遇到问题

为Google Cloud Pub/Sub消息发布方法编写JUnit测试方案

核心思路

直接调用真实的Google Cloud Pub/Sub服务会依赖外部环境,测试成本高且不稳定。我们通过依赖解耦+**Mock框架(Mockito)**模拟Pub/SDK核心对象,实现无外部依赖的单元测试。


步骤1:重构原代码以支持测试

原方法内直接创建Publisher,导致无法Mock。需要抽离Publisher的创建逻辑,用工厂模式解耦:

新增工厂接口与实现

public class PubSubPublisher {
    private final PubSubConfig config;
    private final PublisherFactory publisherFactory;

    // 生产环境构造函数
    public PubSubPublisher(PubSubConfig config) {
        this(config, new DefaultPublisherFactory());
    }

    // 测试用构造函数,注入Mock工厂
    public PubSubPublisher(PubSubConfig config, PublisherFactory publisherFactory) {
        this.config = config;
        this.publisherFactory = publisherFactory;
    }

    // 定义工厂接口
    interface PublisherFactory {
        Publisher createPublisher(TopicName topicName) throws IOException;
    }

    // 默认工厂实现
    static class DefaultPublisherFactory implements PublisherFactory {
        @Override
        public Publisher createPublisher(TopicName topicName) throws IOException {
            return Publisher.newBuilder(topicName).build();
        }
    }

    // 修改原publishJSON方法,使用工厂创建Publisher
    public String publishJSON(String json) throws InterruptedException, IOException, ExecutionException {             
        log.info(" Publishing payload to: "+config.getTopicId());   
        TopicName topicName=TopicName.of(config.getPubsubProjectId(),config.getTopicId());
        Publisher publisher=null;
        try {
            // 替换为工厂创建
            publisher = publisherFactory.createPublisher(topicName);
            ByteString data = ByteString.copyFromUtf8(json);
            PubsubMessage pubsubMessage = PubsubMessage.newBuilder().setData(data).build();
            ApiFuture<String> messageIdFuture = publisher.publish(pubsubMessage);
            String messageId = messageIdFuture.get();
            log.info("Published message ID: " + messageId);
            return messageId;
        } catch (ExecutionException e) {
            log.error("Error while publishing messsage" + e.getMessage());
            throw e;
        } catch (IOException e) {
            log.error( "PubSub exception "+ e.getMessage());
            throw e;
        } catch (InterruptedException e) {
            log.error("Connection making exception for PubSub" + e.getMessage());
            throw e;
        } catch (Exception e) {
            log.error( "publishJSON Error : "+ e.getMessage());
            throw e;    
        } finally {
            if (publisher != null) {
                publisher.shutdown();
                publisher.awaitTermination(1, TimeUnit.MINUTES);
            }
        }
    }
}

步骤2:编写JUnit测试用例(Mockito实现)

测试依赖

确保测试依赖包含:

  • JUnit 5
  • Mockito Core

测试代码示例

import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.InjectMocks;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import com.google.api.core.ApiFuture;
import com.google.cloud.pubsub.v1.Publisher;
import com.google.pubsub.v1.PubsubMessage;
import com.google.pubsub.v1.TopicName;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.*;

@ExtendWith(MockitoExtension.class)
class PubSubPublisherTest {

    @Mock
    private PubSubConfig config;

    @Mock
    private PublisherFactory publisherFactory;

    @Mock
    private Publisher publisher;

    @Mock
    private ApiFuture<String> messageIdFuture;

    @InjectMocks
    private PubSubPublisher pubSubPublisher;

    // 测试正常发布场景
    @Test
    void publishJSON_returnsMessageId_onSuccess() throws Exception {
        // 预设测试数据
        String testProjectId = "demo-project";
        String testTopicId = "demo-topic";
        String testJson = "{\"name\":\"test\"}";
        String expectedMsgId = "msg-12345";

        // 配置Mock行为
        when(config.getPubsubProjectId()).thenReturn(testProjectId);
        when(config.getTopicId()).thenReturn(testTopicId);
        when(publisherFactory.createPublisher(any(TopicName.class))).thenReturn(publisher);
        when(publisher.publish(any(PubsubMessage.class))).thenReturn(messageIdFuture);
        when(messageIdFuture.get()).thenReturn(expectedMsgId);

        // 执行测试方法
        String actualMsgId = pubSubPublisher.publishJSON(testJson);

        // 验证结果
        assert actualMsgId.equals(expectedMsgId);

        // 验证消息内容是否正确
        verify(publisher).publish(argThat(msg -> 
            msg.getData().toStringUtf8().equals(testJson)
        ));

        // 验证资源释放逻辑
        verify(publisher).shutdown();
        verify(publisher).awaitTermination(1, TimeUnit.MINUTES);
    }

    // 测试发布失败抛出异常的场景
    @Test
    void publishJSON_throwsException_onPublishFailure() throws Exception {
        String testJson = "{\"name\":\"test\"}";
        ExecutionException testException = new ExecutionException(new RuntimeException("PubSub service unavailable"));

        // 配置Mock抛出异常
        when(config.getPubsubProjectId()).thenReturn("demo-project");
        when(config.getTopicId()).thenReturn("demo-topic");
        when(publisherFactory.createPublisher(any(TopicName.class))).thenReturn(publisher);
        when(publisher.publish(any(PubsubMessage.class))).thenReturn(messageIdFuture);
        when(messageIdFuture.get()).thenThrow(testException);

        // 验证异常抛出
        try {
            pubSubPublisher.publishJSON(testJson);
            assert false : "Expected ExecutionException was not thrown";
        } catch (ExecutionException e) {
            assert e.getCause().getMessage().equals("PubSub service unavailable");
        }

        // 验证即使异常,资源仍然被释放
        verify(publisher).shutdown();
        verify(publisher).awaitTermination(1, TimeUnit.MINUTES);
    }
}

关键测试点覆盖

  • 正常发布时的消息内容、返回值验证
  • 各类异常场景(ExecutionException、IOException、InterruptedException)的错误抛出验证
  • 无论成功或失败,finally块中的资源释放逻辑是否执行

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 06:10:34