Quarkus中基于SmallRye连接器实现Kafka入站消息拦截器
在Quarkus中使用SmallRye Kafka实现入站消息拦截器(租户切换场景)
完全可以实现Kafka入站消息的拦截,用来读取消息头并切换租户。以下是两种实用的实现方式,适配不同的场景:
方式一:使用原生Kafka ConsumerInterceptor
这种方式和Spring中消费者拦截器的实现逻辑一致,直接在Kafka客户端层面拦截消息:
1. 编写拦截器实现类
import org.apache.kafka.clients.consumer.ConsumerInterceptor; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.OffsetAndMetadata; import org.apache.kafka.common.TopicPartition; import java.util.Map; public class TenantSwitchConsumerInterceptor implements ConsumerInterceptor<String, Object> { @Override public ConsumerRecords<String, Object> onConsume(ConsumerRecords<String, Object> records) { // 遍历消息,读取租户ID并切换上下文 for (ConsumerRecord<String, Object> record : records) { String tenantId = record.headers().lastHeader("tenant-id") != null ? new String(record.headers().lastHeader("tenant-id").value()) : "default"; TenantContext.setCurrentTenant(tenantId); } return records; } @Override public void onCommit(Map<TopicPartition, OffsetAndMetadata> offsets) { // 提交后清空租户上下文,避免线程复用导致的租户污染 TenantContext.clearCurrentTenant(); } @Override public void configure(Map<String, ?> configs) { // 可添加初始化逻辑 } @Override public void close() { // 资源清理逻辑 } }
2. 在配置文件中注册拦截器
在application.properties中添加以下配置,将拦截器绑定到Kafka消费者:
quarkus.kafka.consumer.interceptor.classes=com.yourpackage.TenantSwitchConsumerInterceptor
方式二:使用MicroProfile Reactive Messaging MessageInterceptor
这种方式基于SmallRye Reactive Messaging的拦截机制,更适配Quarkus的反应式编程模型,能处理消息确认、异常等场景:
1. 编写全局消息拦截器
import org.eclipse.microprofile.reactive.messaging.Message; import org.eclipse.microprofile.reactive.messaging.spi.MessageInterceptor; import jakarta.annotation.Priority; import jakarta.enterprise.context.ApplicationScoped; import java.util.Optional; import java.util.concurrent.CompletableFuture; @ApplicationScoped @Priority(100) // 优先级,数字越小越先执行 public class TenantSwitchMessageInterceptor implements MessageInterceptor { @Override public <T> Message<T> intercept(Message<T> message) { // 从Kafka消息元数据中读取租户ID String tenantId = message.getMetadata(org.apache.kafka.common.header.Headers.class) .flatMap(headers -> { org.apache.kafka.common.header.Header header = headers.lastHeader("tenant-id"); return header != null ? Optional.of(new String(header.value())) : Optional.empty(); }) .orElse("default"); // 切换租户上下文 TenantContext.setCurrentTenant(tenantId); // 绑定消息确认/失败后的清理逻辑,确保线程安全 return message.withAck(() -> { TenantContext.clearCurrentTenant(); return CompletableFuture.completedFuture(null); }).withNack(throwable -> { TenantContext.clearCurrentTenant(); return CompletableFuture.completedFuture(null); }); } }
说明
- 这个拦截器会全局作用于所有入站的Kafka消息,无需额外配置。
- 通过
withAck和withNack确保无论消息处理成功还是失败,都会清空租户上下文,避免线程复用带来的问题。
注意事项
- 确保
TenantContext类是线程安全的(例如使用ThreadLocal存储租户ID)。 - 若使用反应式流,需注意上下文传播,避免租户信息丢失。
内容的提问来源于stack exchange,提问作者Leonardo Floriano
相关产品推荐
相关产品推荐

