Kafka Streams单元测试:如何模拟RecordTooLargeException?
解决TopologyTestDriver中模拟RecordTooLargeException的方案
方案1:自定义MockProducer子类,添加消息大小校验
TopologyTestDriver默认使用的MockProducer不会检查消息大小,我们可以继承它并重写send()方法,手动实现大小校验逻辑并抛出异常:
import org.apache.kafka.clients.producer.MockProducer; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.errors.RecordTooLargeException; import java.util.concurrent.CompletableFuture; import java.util.concurrent.Future; public class SizeValidatingMockProducer<K, V> extends MockProducer<K, V> { private final int maxRequestSize; public SizeValidatingMockProducer(int maxRequestSize) { super(true, null, null); this.maxRequestSize = maxRequestSize; } @Override public Future<RecordMetadata> send(ProducerRecord<K, V> record, Callback callback) { int recordSize = calculateApproximateRecordSize(record); if (recordSize > maxRequestSize) { RecordTooLargeException exception = new RecordTooLargeException("Record size exceeds max request limit"); if (callback != null) { callback.onCompletion(null, exception); } return CompletableFuture.failedFuture(exception); } return super.send(record, callback); } private int calculateApproximateRecordSize(ProducerRecord<K, V> record) { // 近似计算消息总大小:键+值的序列化字节数 + Kafka固定元数据开销(约40字节) int keySize = record.key() != null ? serializeKey(record.key()).length : 0; int valueSize = record.value() != null ? serializeValue(record.value()).length : 0; return keySize + valueSize + 40; } }
测试时替换TopologyTestDriver的默认Producer:
int maxRequestSize = 1024; // 配置你的最大请求大小阈值 SizeValidatingMockProducer<String, String> customProducer = new SizeValidatingMockProducer<>(maxRequestSize); TopologyTestDriver testDriver = new TopologyTestDriver(topology, streamsConfig, customProducer);
当发送超过大小限制的消息时,就会触发RecordTooLargeException,进而调用你的自定义ProductionExceptionHandler。
方案2:用Mockito拦截Producer的send方法抛出异常
如果不想自定义MockProducer,可以用Mockito包装原始实例,拦截send方法并在特定条件下抛出异常:
import org.mockito.Mockito; // 创建原始MockProducer MockProducer<String, String> originalProducer = new MockProducer<>(true, null, null); // 用Mockito生成代理实例 MockProducer<String, String> mockedProducer = Mockito.spy(originalProducer); int maxRequestSize = 1024; // 拦截send方法,添加大小校验逻辑 Mockito.doAnswer(invocation -> { ProducerRecord<String, String> record = invocation.getArgument(0); Callback callback = invocation.getArgument(1); int recordSize = calculateApproximateRecordSize(record); if (recordSize > maxRequestSize) { RecordTooLargeException exception = new RecordTooLargeException(); if (callback != null) { callback.onCompletion(null, exception); } return CompletableFuture.failedFuture(exception); } return invocation.callRealMethod(); }).when(mockedProducer).send(Mockito.any(ProducerRecord.class), Mockito.any(Callback.class)); // 传入TopologyTestDriver TopologyTestDriver testDriver = new TopologyTestDriver(topology, streamsConfig, mockedProducer);
方案3:单独测试自定义ProductionExceptionHandler
如果仅需验证Handler的核心逻辑(返回CONTINUE),可以直接针对Handler做单元测试,无需依赖TopologyTestDriver:
import org.junit.jupiter.api.Test; import static org.junit.jupiter.api.Assertions.assertEquals; @Test void testRecordTooLargeExceptionHandling() { CustomProductionExceptionHandler handler = new CustomProductionExceptionHandler(); // 初始化Handler配置 handler.configure(Map.of()); // 构造测试场景 ProducerRecord<String, String> testRecord = new ProducerRecord<>("test-topic", "key", "large-content".repeat(1000)); RecordTooLargeException testException = new RecordTooLargeException(); // 验证返回结果 assertEquals(ProductionExceptionHandlerResponse.CONTINUE, handler.handle(testRecord, testException)); }
内容的提问来源于stack exchange,提问作者BlackViper
相关产品推荐
相关产品推荐

