【发布时间】: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。
我错过了什么?
【问题讨论】:
标签: java spring spring-kafka spring-cloud-stream