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

Kafka Streams中try-catch无法捕获异常的原因咨询

为什么Kafka Streams中的try-catch无法捕获除零异常?

我在使用Spring整合Kafka Streams时,尝试模拟错误场景:当向输入主题传入值0时,会触发除零异常。我将可能引发错误的代码段包裹在try-catch块中,预期异常会被捕获,但实际异常未被捕获,反而抛出如下异常并导致应用停止。

@Component
@Slf4j
public class Topology {

    private static final Serde<String> STRING_SERDE = Serdes.String();

    @Autowired
    void KStreamsTopology(StreamsBuilder streamsBuilder) {
        KStream<String, String> messageStream = streamsBuilder
                .stream(List.of("quickstart-events", "sample-input"), Consumed.with(STRING_SERDE, STRING_SERDE).withName("my-inputs"));
        try {
            messageStream
                    .peek((k, v) -> System.out.println("Input Key: " + k + ", value: " + v))
                    .mapValues(a -> getAnInt(Integer.parseInt(a)))
                    .foreach((k, v) -> System.out.println("output Key: " + k + ", Output value: " + v));
        } catch (Exception exception) {
           log.error("Exceptions occurred in my Topology:: " + exception.getMessage());
        }

    }

    private int getAnInt(Integer value) {
        return 10 / value;
    }
}

抛出的异常信息:

2022-08-10 23:26:27.218  INFO 40343 --- [-StreamThread-1] o.a.k.s.p.internals.StreamThread         : stream-thread [kafka-streams-poc-de132014-fb30-4c10-b5d5-7b190c3e38db-StreamThread-1] Shutdown complete
Exception in thread "kafka-streams-poc-de132014-fb30-4c10-b5d5-7b190c3e38db-StreamThread-1" org.apache.kafka.streams.errors.StreamsException: Exception caught in process. taskId=0_0, processor=my-inputs, topic=sample-input, partition=0, offset=10, stacktrace=java.lang.ArithmeticException: / by zero
    at com.techopact.kafkastreamspoc.topology.Topology.getAnInt(Topology.java:37)
    at com.techopact.kafkastreamspoc.topology.Topology.lambda$KStreamsTopology$1(Topology.java:28)
    at org.apache.kafka.streams.kstream.internals.AbstractStream.lambda$withKey$2(AbstractStream.java:111)
    at org.apache.kafka.streams.kstream.internals.KStreamMapValues$KStreamMapProcessor.process(KStreamMapValues.java:41)
    at org.apache.kafka.streams.processor.internals.ProcessorNode.process(ProcessorNode.java:146)
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forwardInternal(ProcessorContextImpl.java:253)
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:232)
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:191)
    at org.apache.kafka.streams.kstream.internals.KStreamPeek$KStreamPeekProcessor.process(KStreamPeek.java:42)
    at org.apache.kafka.streams.processor.internals.ProcessorNode.process(ProcessorNode.java:146)
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forwardInternal(ProcessorContextImpl.java:253)
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:232)
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:191)
    at org.apache.kafka.streams.processor.internals.SourceNode.process(SourceNode.java:84)
    at org.apache.kafka.streams.processor.internals.StreamTask.lambda$process$1(StreamTask.java:731)
    at org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.maybeMeasureLatency(StreamsMetricsImpl.java:809)
    at org.apache.kafka.streams.processor.internals.StreamTask.process(StreamTask.java:731)
    at org.apache.kafka.streams.processor.internals.TaskManager.process(TaskManager.java:1296)
    at org.apache.kafka.streams.processor.internals.StreamThread.runOnce(StreamThread.java:784)
    at org.apache.kafka.streams.processor.internals.StreamThread.runLoop(StreamThread.java:604)
    at org.apache.kafka.streams.processor.internals.StreamThread.run(StreamThread.java:576)

    at org.apache.kafka.streams.processor.internals.StreamTask.process(StreamTask.java:758)
    at org.apache.kafka.streams.processor.internals.TaskManager.process(TaskManager.java:1296)
    at org.apache.kafka.streams.processor.internals.StreamThread.runOnce(StreamThread.java:784)
    at org.apache.kafka.streams.processor.internals.StreamThread.runLoop(StreamThread.java:604)
    at org.apache.kafka.streams.processor.internals.StreamThread.run(StreamThread.java:576)
Caused by: java.lang.ArithmeticException: / by zero
    at com.techopact.kafkastreamspoc.topology.Topology.getAnInt(Topology.java:37)
    at com.techopact.kafkastreamspoc.topology.Topology.lambda$KStreamsTopology$1(Topology.java:28)
    at org.apache.kafka.streams.kstream.internals.AbstractStream.lambda$withKey$2(AbstractStream.java:111)
    at org.apache.kafka.streams.kstream.internals.KStreamMapValues$KStreamMapProcessor.process(KStreamMapValues.java:41)
    at org.apache.kafka.streams.processor.internals.ProcessorNode.process(ProcessorNode.java:146)
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forwardInternal(ProcessorContextImpl.java:253)
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:232)
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:191)
    at org.apache.kafka.streams.kstream.internals.KStreamPeek$KStreamPeekProcessor.process(KStreamPeek.java:42)
    at org.apache.kafka.streams.processor.internals.ProcessorNode.process(ProcessorNode.java:146)
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forwardInternal(ProcessorContextImpl.java:253)
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:232)
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:191)
    at org.apache.kafka.streams.processor.internals.SourceNode.process(SourceNode.java:84)
    at org.apache.kafka.streams.processor.internals.StreamTask.lambda$process$1(StreamTask.java:731)
    at org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.maybeMeasureLatency(StreamsMetricsImpl.java:809)
    at org.apache.kafka.streams.processor.internals.StreamTask.process(StreamTask.java:731)
    ... 4 more
2022-08-10 23:31:26.921  INFO 40343 --- [90c3e38db-admin] org.apache.kafka.clients.NetworkClient   : [AdminClient clientId=kafka-streams-poc-de132014-fb30-4c10-b5d5-7b190c3e38db-admin] Node -1 disconnected.
2022-08-10 23:36:27.007  INFO 40343 --- [90c3e38db-admin] org.apache.kafka.clients.NetworkClient   : [AdminClient clientId=kafka-streams-poc-de132014-fb30-4c10-b5d5-7b190c3e38db-admin] Node 0 disconnected.

我已知Kafka Streams应用的正确异常处理方式,但仍想了解为何此异常未被try-catch块捕获。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 02:15:37