Kafka 实战
从环境搭建、消息系统设计到生产消费、重试、CDC 与消息兼容性。
7 篇文章 · 按阅读顺序排列
Kafka(一)使用Docker Compose安装单机Kafka+Kafka UI+Prometheus JMX Exporter
使用Docker Compose搭建KRaft模式的单机Kafka,并集成Kafka UI和Prometheus JMX Exporter。说明监听器、角色、内存和监控配置,记录启动验证、JMX端口冲突及资源配置问题的排查方法。
Kafka(二)将邮件发送从业务系统中解耦之消息系统设计
以邮件发送为例,介绍通过Kafka解耦业务系统与邮件系统的设计思路。讨论两者的职责边界、统一消息格式、发送与回调流程,以及新旧系统切换时的回退方案,为后续生产者和消费者实现提供背景。
Kafka(三)生产者发送JSON消息+使用统一序列化器+提升吞吐量
介绍Kafka生产者发送JSON消息的实现,使用统一序列化器处理不同对象,并配置异步发送、失败记录和重试。结合邮件消息的实际大小,讨论压缩、批次大小、请求大小和等待时间对吞吐量的影响。
Kafka(四)消费者消费JSON消息+使用统一反序列化器+提升吞吐量
介绍Kafka消费者处理JSON消息的实现,包括统一反序列化、消息去重、失败重试与偏移量提交。结合压测讨论分区数和消费参数的调整,并使用模板方法复用不同优先级邮件的消费流程。
Kafka(五)消费者回调 +定时重试 + 理解Rebalance
至此为止,利用Kafka实现一个消息系统就基本完成了,所有关键的代码都在不同的博文中并进行了详细说明,说过想要体会完整的设计、实现思路,请移步源码仓库获取完整代码。下一篇关于Kafak的博文打算分享一下如何利用Kafka Connect将Oracle数据库的数据同步到Postgre SQL中。
Kafka(六)利用Kafka Connect+Debezium通过CDC方式将Oracle数据库的数据同步至PostgreSQL中以及实现缓存一致性
介绍CDC的基本思路,并通过Kafka Connect和Debezium示例,将Oracle数据同步到PostgreSQL。讨论数据库同步、Elasticsearch数据导入及缓存一致性等场景,并比较不同方案的适用条件。
Kafka(七)集成Apache Avro+Apicurio Schema Registry以保障生产者与消费者的消息兼容性
针对Kafka消息格式变更引发的反序列化失败,介绍从JSON迁移到Avro并接入Apicurio Schema Registry的流程。涵盖Schema定义、Java类生成、生产与消费配置,以及版本兼容性检查,降低毒药消息带来的风险。