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;
}
}
延伸阅读
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。