08. Java Spring Kafka 实战

Spring Kafka 配置、@KafkaListener、并发消费者、事务与错误处理

1. 配置

spring:
  kafka:
    bootstrap-servers: kafka1:9092,kafka2:9092
    producer:
      acks: all
      retries: 3
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.apache.kafka.common.serialization.StringSerializer
    consumer:
      group-id: order-service
      auto-offset-reset: earliest
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer

2. 生产消息

@Service
public class OrderEventPublisher {
    @Autowired private KafkaTemplate<String, String> kafkaTemplate;

    public void publishOrderCreated(Order order) {
        kafkaTemplate.send("orders.created", order.getUserId(),
            JsonUtils.toJson(order));
    }
}

3. 消费消息

@Component
public class OrderEventConsumer {

    @KafkaListener(topics = "orders.created",
                   groupId = "order-service",
                   concurrency = "3")  // 3 个并发消费者
    public void onOrderCreated(ConsumerRecord<String, String> record,
                               Acknowledgment ack) {
        Order order = JsonUtils.fromJson(record.value(), Order.class);
        processOrder(order);
        ack.acknowledge();  // 手动提交
    }

    @KafkaListener(topics = "orders.cancelled")
    public void onOrderCancelled(String message) {
        // 处理取消
    }
}

4. 错误处理

@Configuration
public class KafkaConfig {
    @Bean
    public DefaultErrorHandler errorHandler() {
        BackOff fixedBackOff = new FixedBackOff(2000, 3);  // 重试3次,间隔2s
        DefaultErrorHandler handler = new DefaultErrorHandler(fixedBackOff);
        handler.addNotRetryableExceptions(IllegalArgumentException.class);
        return handler;
    }
}

延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「kafka」更多文章

  1. 事件驱动架构:Event Sourcing、CQRS 与 Saga 模式
  2. Kafka 运维监控与故障恢复:JMX 指标、Lag 监控与分区重分配
  3. Kafka 详解:分布式日志系统、ISR 与一致性保证