【问题标题】:How to configure GlobalChannelInterceptor for spring-cloud-stream?如何为 spring-cloud-stream 配置 GlobalChannelInterceptor?
【发布时间】:2019-05-31 16:09:45
【问题描述】:

我有一个使用 kafka 绑定的 spring-cloud-stream worker

@Slf4j
@EnableBinding(KafkaStreamsProcessor.class)
@RequiredArgsConstructor
public class SomeWorker {

    @StreamListener(Sink.INPUT)
    @SendTo(Source.OUTPUT)
    public KStream<?, Obj> process(KStream<?, Obj> objStream) {
        return objStream.something();
    }
}

还有一个全局拦截器

@Component
@Slf4j
@GlobalChannelInterceptor
public class StreamInterceptor implements ChannelInterceptor {

    @Override
    public Message<?> preSend(Message<?> msg, MessageChannel mc) {
        log.info("In preSend");
        return msg;
    }

    @Override
    public void postSend(Message<?> msg, MessageChannel mc, boolean bln) {
        log.info("In postSend");
    }

    @Override
    public void afterSendCompletion(Message<?> msg, MessageChannel mc, boolean bln, Exception excptn) {
        log.info("In afterSendCompletion");
    }

    @Override
    public boolean preReceive(MessageChannel mc) {
        log.info("In preReceive");
        return true;
    }

    @Override
    public Message<?> postReceive(Message<?> msg, MessageChannel mc) {
        log.info("In postReceive");
        return msg;
    }

}

在接收到任何流式消息时,不会调用 GlobalChannelInterceptor。

我错过了什么?

ps:我正在关注这个测试https://github.com/spring-cloud/spring-cloud-stream/blob/master/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/interceptor/BoundChannelsInterceptedTest.java#L67

【问题讨论】:

    标签: java spring spring-kafka spring-cloud-stream


    【解决方案1】:

    Kafka Streams binder 不是基于MessageChannels,所以没有通道可以拦截。

    【讨论】:

    • @Garry Russel,有什么办法可以拦截像ChannelInterceptor这样的消息?
    • 我认为唯一的方法是实现一个Transformer并将其添加到流拓扑中。
    猜你喜欢
    • 2018-05-27
    • 2021-04-03
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-10-24
    • 2020-01-28
    • 2020-04-07
    相关资源
    最近更新 更多