ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

【基于 Swoole+Hyperf 的微服务实战】第七周·周三:Kafka 与高吞吐消息场景

【基于 Swoole+Hyperf 的微服务实战】第七周·周三:Kafka 与高吞吐消息场景 【基于 SwooleHyperf 的微服务实战】第七周·周三Kafka 与高吞吐消息场景今天我们进入第七周周三主题是Kafka 与高吞吐消息场景。前两天我们基于 RabbitMQ 构建了可靠的消息通信但 RabbitMQ 在处理海量日志、流式数据时可能会遇到吞吐瓶颈。今天我们将引入Apache Kafka——一个分布式、高吞吐、支持持久化的流处理平台。我们将使用hyperf/kafka组件实现高并发日志收集将用户行为日志写入 Kafka并由消费者实时处理体验 Kafka 在海量数据下的强劲性能。今日目标理解 Kafka 的核心概念Topic、Partition、Broker、Consumer Group、Offset。使用 Docker 部署 KafkaKRaft 模式无需 ZooKeeper。安装hyperf/kafka配置生产者和消费者。实现一个用户行为日志收集场景HTTP 请求日志通过 Kafka 异步发送消费者批量写入数据库或对象存储。通过压测对比 Kafka 与 RabbitMQ 的吞吐能力感受 Kafka 在流式数据处理中的优势。一、环境准备部署 Kafka约 30 分钟1. 添加 Kafka 服务到 Docker ComposeKafka 3.x 支持KRaft模式无 ZooKeeper。编辑swoole-course/docker-compose.yml增加kafka:image:confluentinc/cp-kafka:7.5.0container_name:kafka-labports:-9092:9092-9093:9093# 内部监听可忽略environment:KAFKA_NODE_ID:1KAFKA_LISTENER_SECURITY_PROTOCOL_MAP:CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXTKAFKA_LISTENERS:PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093KAFKA_ADVERTISED_LISTENERS:PLAINTEXT://kafka:9092KAFKA_PROCESS_ROLES:broker,controllerKAFKA_CONTROLLER_QUORUM_VOTERS:1kafka:9093KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR:1KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR:1KAFKA_TRANSACTION_STATE_LOG_MIN_ISR:1restart:unless-stopped启动 Kafkadocker-composeup-dkafka确认 Kafka 启动成功可能需要几十秒dockerlogs kafka-lab|grepstarted2. 安装 hyperf/kafka 组件进入 PHP 容器docker-composeexecswoolebashcd/var/www/hyperf-appcomposerrequire hyperf/kafka发布配置php bin/hyperf.php vendor:publish hyperf/kafka这将在config/autoload/kafka.php生成默认配置。二、知识核心Kafka 架构与高吞吐原理约 1 小时1. Kafka 基本概念Topic消息分类类似 RabbitMQ 的 routing key 和队列的组合。每个 Topic 可以有多个分区。Partition一个 Topic 被分割为多个分区每个分区是一个有序的、不可变的消息序列。分区是 Kafka 并行处理的核心可以分布在不同的 Broker 上。BrokerKafka 服务节点。一个集群由多个 Broker 组成每个 Broker 承担一部分分区。Producer发布消息到指定 Topic可以指定分区键或由 Kafka 自动分区。Consumer Group消费者组组内多个消费者协同消费一个 Topic 的多个分区每个分区只能被组内一个消费者消费。水平扩展消费者数量时分区数要与之匹配。Offset消费者在分区内的位置标记用于记录已消费位置支持回溯和重放。2. 与 RabbitMQ 的对比特性RabbitMQKafka设计目标可靠消息传递灵活路由高吞吐、流式数据处理消息模型交换机队列支持复杂路由发布-订阅Topic/Partition吞吐量万级 QPS百万级 QPS持久化内存磁盘默认可靠默认磁盘顺序写入极快消息回溯不支持消费后删除支持可按 Offset 重放历史数据事务支持支持幂等生产者事务适用场景业务命令、RPC 回调、延迟消息日志收集、实时流计算、事件溯源今天我们的场景是用户行为日志收集非常适合 Kafka因为日志量大、允许少量延迟、需要持久化和顺序保证。3. hyperf/kafka 组件基于longlang/phpkafkaSwoole 协程化的 Kafka 客户端提供协程化的生产者支持批量发送。消费者可以按分区消费并手动提交 Offset。支持 SASL 认证、SSL 加密等。三、实战构建用户行为日志 Kafka 管道约 2.5 小时步骤 1配置 Kafka 连接编辑config/autoload/kafka.php?phpreturn[default[hostenv(KAFKA_HOST,kafka),port(int)env(KAFKA_PORT,9092),pool[min_connections1,max_connections10,connect_timeout10.0,wait_timeout3.0,heartbeat-1,max_idle_time60,],],];步骤 2创建日志生产者新建app/Kafka/Producer/LogProducer.php?phpnamespaceApp\Kafka\Producer;useHyperf\Kafka\AbstractProducer;useHyperf\Kafka\Annotation\Producer;#[Producer(topic:user_behavior_logs,name:LogProducer)]classLogProducerextendsAbstractProducer{// 无需额外代码基类已提供 send 方法}说明AbstractProducer提供了send(string $payload, string $key null)方法$key用于分区路由相同 key 进入同一分区保证顺序。步骤 3在业务中发送日志消息我们可以在网关的 JWT 中间件或控制器中记录用户每次请求的行为日志。为了演示在hyperf-app的ArticleController::show()方法中发送一条日志消息。修改app/Controller/ArticleController.phpuseApp\Kafka\Producer\LogProducer;useHyperf\Di\Annotation\Inject;#[Inject]privateLogProducer$logProducer;publicfunctionshow(int$id){// 原有业务逻辑...$articleArticle::with([category,tags])-find($id);if(!$article){/* ... */}// 发送用户行为日志到 Kafka$userId$this-request-getHeaderLine(X-User-Id)??0;$logDatajson_encode([user_id(int)$userId,actionview_article,article_id$id,timestamptime(),ip$this-request-getServerParams()[remote_addr]??,]);$this-logProducer-send($logData,user_.$userId);// 以用户ID为分区键同一用户日志有序// 返回文章...}这样每次查看文章详情时一条日志就会发送到 Kafka Topicuser_behavior_logs。步骤 4创建日志消费者新建app/Kafka/Consumer/LogConsumer.php?phpnamespaceApp\Kafka\Consumer;useHyperf\Kafka\AbstractConsumer;useHyperf\Kafka\Annotation\Consumer;useHyperf\Kafka\Result;uselonglang\phpkafka\Consumer\ConsumeMessage;#[Consumer(topic:user_behavior_logs,groupId:log-processor,name:LogConsumer,nums:2)]classLogConsumerextendsAbstractConsumer{publicfunctionconsume(ConsumeMessage$message):string{$datajson_decode($message-getValue(),true);if(!$data){returnResult::ACK;}$userId$data[user_id]??0;$action$data[action]??;echo[Kafka消费者] 处理日志: 用户{$userId}{$action}文章{$data[article_id]}\n;// 模拟批量写入数据库或 ES这里仅打印// 实际可写入 MySQL 日志表或发送到 Elasticsearch 用于分析returnResult::ACK;}}要点groupId消费者组名称同一组的消费者分配不同分区保证每条消息只被组内一个消费者处理。nums启动的消费者协程数应小于等于 Topic 的分区数才能充分利用并行性。ConsumeMessage提供getValue()、getKey()、getOffset()等方法。步骤 5配置消费者进程编辑config/autoload/processes.php添加 Kafka 消费者进程如果没有则创建文件?phpreturn[\Hyperf\Kafka\Process\ConsumerProcess::class,];重启hyperf-appphp bin/hyperf.php start消费者进程会启动自动加入log-processor消费者组消费user_behavior_logs分区。步骤 6创建 Topic手动或自动Kafka 默认允许自动创建 Topicauto.create.topics.enabletrue但建议手动创建以控制分区数。在 Kafka 容器中执行dockerexec-itkafka-lab kafka-topics--create--topicuser_behavior_logs --bootstrap-server kafka:9092--partitions3--replication-factor1查看 Topicdockerexec-itkafka-lab kafka-topics--list--bootstrap-server kafka:9092四、成果测试与高吞吐验证约 1 小时1. 功能验证请求文章详情接口curlhttp://localhost:9501/articles/1观察控制台消费者日志输出确认日志消息被消费。查看 Kafka 消息使用控制台消费者dockerexec-itkafka-lab kafka-console-consumer--topicuser_behavior_logs --bootstrap-server kafka:9092 --from-beginning能看到历史消息。2. 压测吞吐量对比使用wrk或ab对文章详情接口进行高并发压测同时观察 Kafka 生产者和消费者的处理速度。# 压测 20000 请求并发 100wrk-t4-c100-d60shttp://localhost:9501/articles/1RabbitMQ场景下昨天消息吞吐受单队列限制压测时可能堆积。Kafka场景下多分区并行消费者组中多个协程并发消费吞吐极高。可以在 Kafka 管理工具如kafka-consumer-groups查看消费延迟dockerexec-itkafka-lab kafka-consumer-groups --bootstrap-server kafka:9092--grouplog-processor--describe观察LAG未消费消息数是否迅速降为零。3. 测试消费者组和分区重平衡启动两个hyperf-app实例不同 Worker 或多个容器它们属于同一消费者组观察 Kafka 自动将分区均匀分配。停止其中一个实例剩余实例会自动接管全部分区Rebalance。4. 测试清单检验项方法通过标准生产者发送消息查看 Kafka 控制台消费者或 Topic 消息数消息成功进入 Topic消费者接收消息查看应用日志输出打印出日志信息且LAG近 0分区键路由相同user_id的消息Kafka 工具查看分区进入同一分区消费者组并行高并发下消费速度跟得上生产速度无明显堆积Offset 自动提交重启消费者后不会重复消费已处理的消息消息数量不重复需手动确认配置高吞吐稳定性长时间压测系统无 OOM连接正常持续稳定常见问题Topic 不存在生产者发送时会自动创建如果开启但建议手动创建指定分区数。消费者不消费检查 Group ID 是否一致Topic 名称拼写以及分区数是否大于消费者nums。连接 Kafka 失败确认kafka主机名可达PHP 容器内 ping kafka或使用host.docker.internal。五、今日作业与学习产出提交代码将LogProducer、LogConsumer、Kafka 配置文件、进程配置等提交。完善日志场景实现一个日志持久化消费者将行为日志批量写入 MySQL使用事务或写入 Elasticsearch 供 Kibana 可视化。添加异常行为检测消费者分析日志若某用户 1 分钟内请求超过 100 次发送告警消息到 RabbitMQ 或直接钉钉。学习笔记画出 Kafka 的分区、消费者组模型图说明如何水平扩展消费者。对比 RabbitMQ 与 Kafka 的使用场景明确在架构中何时选择哪一个。挑战任务使用Kafka Streams或KSQL实现对日志的实时聚合统计如热门文章 TopN并暴露 API 供查询。配置Kafka Connect将日志直接写入 S3 对象存储实现数据湖入湖。通过今天的学习你已掌握高吞吐场景下的异步消息利器 Kafka。现在你的微服务同时拥有了 RabbitMQ可靠命令和 Kafka流式日志两种消息引擎可以根据业务特点灵活选择。明天我们将结合事件驱动设计实现内部服务的完全解耦与最终一致性。
返回列表