ARTICLE DETAIL

资讯详情

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

AutoMQ 中的 Trogdor 混沌测试框架:Kafka 基准测试与故障注入实战指南

AutoMQ 中的 Trogdor 混沌测试框架:Kafka 基准测试与故障注入实战指南 AutoMQ 中的 Trogdor 混沌测试框架Kafka 基准测试与故障注入实战指南【免费下载链接】automqDiskless Kafka® on S3. 10x Cost-Effective. No Cross-AZ Traffic Cost. Autoscale in seconds. Single-digit ms latency. Multi-AZ Availability.项目地址: https://gitcode.com/GitHub_Trending/au/automq本文以 trogdor/README.md 为核心骨架结合本仓库trogdor/、tests/spec/与config/下的源码与配置文件完整讲解 Trogdor 的架构原理、任务模型、内置工作负载与故障注入能力以及从单机 Quickstart 到 Exec 模式的完整实操流程。读完本文你将掌握如何用 Trogdor 对 Kafka含 AutoMQ 的 diskless 架构进行压测、延迟测量与混沌演练并能读懂并自定义任务 Spec。什么是 TrogdorTrogdor 是 Apache Kafka 生态中的一套测试框架test framework它面向运行中的 Kafka 集群提供两类核心能力运行基准测试与各种工作负载workloads例如生产者压测、消费者压测、生产消费回环测试并统计延迟分位数注入故障faults例如人为停住某个进程、制造网络分区从而对系统进行压力测试chaos testing验证集群在异常场景下的行为。在本仓库AutoMQ中Trogdor 位于 trogdor/ 目录下其 Java 包结构org.apache.kafka.trogdor.*与 Kafka 上游保持一致因此既可用于传统 Kafka 集群也可用于 AutoMQ 基于 S3 的 diskless Kafka 集群做压测与故障演练。配套的示例 Spec 位于 tests/spec/配置文件模板位于 config/trogdor.conf。快速开始单节点集群跑通第一个压测任务第一步启动 ZooKeeper 与 Kafka BrokerTrogdor 需要一个可用的 Kafka 集群作为被测对象。先在单节点上启动 ZooKeeper./bin/zookeeper-server-start.sh ./config/zookeeper.properties /tmp/zookeeper.log 再启动 Kafka Broker./bin/kafka-server-start.sh ./config/server.properties /tmp/kafka.log 说明AutoMQ 支持 KRaft 模式对应配置可参考 config/kraft/ 下的broker.properties、controller.properties与server.properties。Trogdor 只通过bootstrapServers与集群交互并不关心元数据模式。第二步启动 Trogdor Agent 与 CoordinatorTrogdor 由两个守护进程组成Agent负责在节点上执行任务和Coordinator负责任务编排。启动命令如下# 启动 Agent ./bin/trogdor.sh agent -c ./config/trogdor.conf -n node0 /tmp/trogdor-agent.log # 启动 Coordinator ./bin/trogdor.sh coordinator -c ./config/trogdor.conf -n node0 /tmp/trogdor-coordinator.log 两个命令都用-c指定配置文件、用-n指定节点名。默认配置文件 config/trogdor.conf 内容如下{ platform: org.apache.kafka.trogdor.basic.BasicPlatform, nodes: { node0: { hostname: localhost, trogdor.agent.port: 8888, trogdor.coordinator.port: 8889 } } }该配置声明了平台实现类为BasicPlatform并定义了节点node0hostname 为localhostAgent 默认监听 8888 端口Coordinator 默认监听 8889 端口。这些默认值与源码中Agent.DEFAULT_PORT 8888、Coordinator.DEFAULT_PORT 8889一致见 Agent.java 与 Coordinator.java。第三步确认所有守护进程已启动jps 116212 Coordinator 115188 QuorumPeerMain 116571 Jps 115420 Kafka 115694 Agent看到Coordinator、QuorumPeerMainZooKeeper、KafkaBroker与Agent四个进程即表示环境就绪。第四步提交一个压测任务使用 Trogdor 自带命令行客户端提交任务./bin/trogdor.sh client createTask -t localhost:8889 -i produce0 --spec ./tests/spec/simple_produce_bench.json Sent CreateTaskRequest for task produce0.-t指定 Coordinator 地址localhost:8889-i指定任务 IDproduce0--spec指向任务 Spec 文件。这里用到的示例 Spec 位于 tests/spec/simple_produce_bench.json{ class: org.apache.kafka.trogdor.workload.ProduceBenchSpec, durationMs: 10000000, producerNode: node0, bootstrapServers: localhost:9092, targetMessagesPerSec: 10000, maxMessages: 50000, activeTopics: { foo[1-3]: { numPartitions: 10, replicationFactor: 1 } }, inactiveTopics: { foo[4-5]: { numPartitions: 10, replicationFactor: 1 } } }第五步查询任务状态./bin/trogdor.sh client showTask -t localhost:8889 -i produce0 Task bar of type org.apache.kafka.trogdor.workload.ProduceBenchSpec is DONE. FINISHED at 2019-01-09T20:38:22.039-08:00 after 6s任务状态为DONE表示已完成。第六步查看压测结果./bin/trogdor.sh client showTask -t localhost:8889 -i produce0 --show-status Task bar of type org.apache.kafka.trogdor.workload.ProduceBenchSpec is DONE. FINISHED at 2019-01-09T20:38:22.039-08:00 after 6s Status: { totalSent : 50000, averageLatencyMs : 17.83388, p50LatencyMs : 12, p95LatencyMs : 75, p99LatencyMs : 96, transactionsCommitted : 0 }输出中的totalSent共发送 50000 条、averageLatencyMs平均延迟以及p50/p95/p99LatencyMs中位数、95 分位、99 分位延迟正是判断集群吞吐与延迟特性的核心指标。Trogdor 架构Coordinator 与 Agent 的分工Trogdor 采用单 Coordinator 管理多 Agent的星型架构Coordinator协调者整个集群只有一个。它负责管理任务task每个任务可以是基准测试、故障注入或工作负载。为了执行任务Coordinator 会在一个或多个 Agent 节点上创建 worker。Coordinator 的核心逻辑在 Coordinator.java 与 TaskManager.java 中。Agent代理每个集群节点对应一个 Agent 进程负责真正执行任务。例如运行工作负载时实际生产/消费消息的进程就是 Agent。Agent 由 Agent.java 与 WorkerManager.java 实现。通信层面Coordinator 与 Agent 都暴露基于 JSON 序列化对象的 REST 接口通过 JsonRestServer.java 实现。命令行程序trogdor.sh client的存在就是为了免去手工构造 JSON 消息体即可向两者发送请求。一个值得注意的协议特性除 shutdown 请求外所有 Trogdor RPC 都是幂等的——连续发送两次相同的 RPC效果与发送一次完全相同。这让客户端在重试时不用担心副作用。任务模型Spec、状态机与不可变性任务 Spec 的通用字段每个任务由一份**规格说明specificationSpec**描述其中包含class任务类型值为完整的 Java 类名如org.apache.kafka.trogdor.workload.ProduceBenchSpecstartMs任务开始时间单位为自 UNIX epoch 起的毫秒数durationMs任务持续时间单位为毫秒其他任务特有字段由具体 Spec 定义。在源码层面所有 Spec 都继承自抽象基类 TaskSpec.java。该类通过 Jackson 的JsonTypeInfo(use JsonTypeInfo.Id.CLASS, property class)实现按class字段反序列化多态 Spec并通过startMs/durationMs构造器参数完成通用字段的解析。其中durationMs会被裁剪到[0, MAX_TASK_DURATION_MS]区间MAX_TASK_DURATION_MS 1000000000000000L用于避免 64 位溢出与浮点舍入问题——因为 JSON 规范中数值本质上是浮点数。任务 Spec 是不可变的任务创建之后其描述不会发生任何变化。这保证了同一个 Spec 在不同节点、不同时间点上语义一致也简化了状态同步。任务状态机任务在生命周期中会经历以下几个状态状态含义PENDING任务等待执行RUNNING任务正在运行STOPPING任务正在停止过程中DONE任务已完成处于DONE状态的任务还带有一个error字段——如果任务失败该字段会被设置。一个网络分区任务示例下面的 Spec 描述了一次在节点 1 和节点 2、3 之间的网络分区即node1、node2组成一个分区node3单独一个分区持续 30 秒{ class: org.apache.kafka.trogdor.fault.NetworkPartitionFaultSpec, startMs: 1000, durationMs: 30000, partitions: [[node1, node2], [node3]] }一个完整 ProduceBench 任务示例下面这个任务在单生产者节点上运行 ProduceBench创建 5 个 topic、以每秒 10000 条的速率生产。键按顺序生成并使用配置的分区器DefaultPartitioner{ class: org.apache.kafka.trogdor.workload.ProduceBenchSpec, durationMs: 10000000, producerNode: node0, bootstrapServers: localhost:9092, targetMessagesPerSec: 10000, maxMessages: 50000, activeTopics: { foo[1-3]: { numPartitions: 10, replicationFactor: 1 } }, inactiveTopics: { foo[4-5]: { numPartitions: 10, replicationFactor: 1 } }, keyGenerator: { type: sequential, size: 8, offset: 1 }, useConfiguredPartitioner: true }任务提交给 Coordinator 后Coordinator 判断到达startMs便在相应 Agent 上创建 workerworker 一直运行直到任务结束。Topic 通配语法上述示例中的foo[1-3]是一种范围展开语法它会被展开为foo1、foo2、foo3三个 topic。该逻辑由 StringExpander.java 实现TopicsSpec见 TopicsSpec.java在反序列化时即完成展开并生成不可变副本保证 Spec 不可变性。内置工作负载Workloads工作负载workload在集群上执行操作并测量性能当操作无法执行时工作负载会失败。Trogdor 内置了以下三种工作负载。ProduceBench生产者压测ProduceBench 在单个 Agent 节点上启动一个 Kafka 生产者向若干分区生产消息并测量平均生产延迟以及中位数、95 分位、99 分位延迟。其 Spec 类为 ProduceBenchSpec.javaWorker 实现为 ProduceBenchWorker.java。ProduceBench 的关键配置字段均可在 Spec JSON 中设置字段含义默认行为producerNode运行生产者的节点名空串bootstrapServersKafka broker 地址列表空串targetMessagesPerSec目标生产速率条/秒无maxMessages最多生产的消息数无keyGenerator键生成器如sequential顺序生成默认SequentialPayloadGenerator(4, 0)valueGenerator值生成器默认 512 字节常量负载transactionGenerator事务生成器配合事务性生产者按时间间隔或消息数量提交事务无非事务producerConf/commonClientConf/adminClientConf生产者/公共客户端/AdminClient 的额外配置MapString,String空 MapactiveTopics生产消息的目标 topic 集合空inactiveTopics只创建、不生产的 topic 集合空useConfiguredPartitioner是否使用配置的分区器true则用 DefaultPartitionerfalseskipFlush是否跳过 flush 直接统计false从源码看ProduceBenchSpec的newController返回的 TaskController 会把任务限定到producerNode单一节点上执行topology - Collections.singleton(producerNode)也就是说 ProduceBench 的负载完全由该 Agent 承担。事务支持通过TransactionGenerator接口实现仓库提供了基于时间间隔与消息数量的两种实现TimeIntervalTransactionsGenerator与UniformTransactionsGenerator。RoundTripWorkload生产消费回环RoundTripWorkload 同时测试生产与消费它在单个节点上启动一个 Kafka 生产者和一个消费者消费者读回生产者生产的消息从而验证消息投递的全链路正确性。其 Spec 为 RoundTripWorkloadSpec.javaWorker 为 RoundTripWorker.java。RoundTripWorkload 的字段与 ProduceBench 类似但多出consumerConf消费者配置并且默认valueGenerator为 32 字节的均匀随机负载UniformRandomPayloadGenerator(32, 123, 10)。该工作负载适合验证生产后能完整消费回来的端到端数据一致性常用于数据完整性回归测试。ConsumeBench消费者压测ConsumeBench 在单个 Agent 节点上启动一个或多个 Kafka 消费者并根据配置见 ConsumeBenchSpec.java选择两种消费模式之一subscribe 模式消费者通过 consumer group 功能订阅一组 topic由 group 动态分配分区assign 模式消费者手动将分区分配给自己activeTopics中包含topic:partition形式的条目时触发。ConsumeBench 测量平均消费延迟以及中位数、95 分位、99 分位延迟。ConsumeBenchSpec 的activeTopics字段支持范围展开例如foo[1-3]:[1-3]会展开为foo1:1, foo1:2, foo1:3, foo2:1, ..., foo3:3的全部 9 个分区组合。其语义规则如下均可从源码注释与materializeTopics()方法确认只要存在至少一个topic:partition对消费者就使用KafkaConsumer.assign()手动分配此时若某 topic 未指定分区如[foo:1, bar]中的bar则消费者会抓取并分配该 topic 的全部分区若没有任何topic:partition对消费者使用KafkaConsumer.subscribe()订阅 topic由 consumer group 动态分配分区threadsPerWorker控制单个 worker 内启动的消费者线程数默认值为 1targetMessagesPerSec、maxMessages、activeTopics均对每个消费者单独生效未指定consumerGroup时每个消费者分配一个不同的随机 group指定时所有消费者使用同一个 group。由于同一 group 内两个消费者不能分配到同一分区当threadsPerWorker 1且显式指定分区时会抛出ConfigException中止任务recordProcessor字段允许在消息 poll 后立即执行自定义处理逻辑接口见 RecordProcessor.java。示例消费者属于 groupcg订阅foo1、foo2、foo3与bar四个 topic{ class: org.apache.kafka.trogdor.workload.ConsumeBenchSpec, durationMs: 10000000, consumerNode: node0, bootstrapServers: localhost:9092, maxMessages: 100, consumerGroup: cg, activeTopics: [foo[1-3], bar] }故障注入Faults故障注入会故意破坏集群中的某些东西用于验证系统在异常情况下的健壮性。Trogdor 内置了两种经典故障。ProcessStopFault进程冻结ProcessStopFault 通过向进程发送SIGSTOP信号使其停止当故障结束时再以SIGCONT信号恢复进程。该故障很适合模拟节点无响应但未宕机的卡顿场景。其 Spec 与 Worker 分别为 ProcessStopFaultSpec.java 与 ProcessStopFaultWorker.java。NetworkPartitionFault网络分区NetworkPartitionFault 在一组或多组节点之间制造人为的网络分区。当前实现基于iptables规则设置在受影响节点的出站流量上因此受影响节点仍然可以从集群外部被访问到。也就是说被隔离的节点仍然可以接收来自外部的连接但无法向分区内的其他节点正常发包。其 Spec 为 NetworkPartitionFaultSpec.java。从源码的partitionSets()方法可以推断出约束一个节点不能出现在多个分区中否则会抛出RuntimeException。partitions字段是一个二维数组每个子数组代表一个岛。例如[[node1, node2], [node3]]表示node1与node2相互可达而node3与两者隔离。相关控制逻辑位于 NetworkPartitionFaultController.java 与 NetworkPartitionFaultWorker.java。补充故障包 trogdor/src/main/java/org/apache/kafka/trogdor/fault/ 下还包含DegradedNetworkFaultSpec/Worker网络劣化、FilesUnreadableFaultSpec、KiboshFaultWorker等更多故障类型需要更细粒度故障注入时可进一步研读。外部进程支持ExternalCommandWorkerTrogdor 支持在外部进程中运行任意命令——这是一种通用的扩展机制可以把任何可配置命令放进 Trogdor 框架中执行无论它是 Python 程序、bash 脚本、Docker 镜像还是其他可执行体。ExternalCommandSpecExternalCommandSpec.java 描述这类任务核心字段为command要执行的命令及其参数列表和workload传给外部进程的工作负载 JSON 对象。示例{ class: org.apache.kafka.trogdor.workload.ExternalCommandSpec, command: [python, /path/to/trogdor/python/runner], durationMs: 10000000, commandNode: node0, workload: { class: org.apache.kafka.trogdor.workload.ProduceBenchSpec, bootstrapServers: localhost:9092, targetMessagesPerSec: 10, maxMessages: 100, activeTopics: { foo[1-3]: { numPartitions: 3, replicationFactor: 1 } }, inactiveTopics: { foo[4-5]: { numPartitions: 3, replicationFactor: 1 } } } }Worker 与外部进程的 JSON 通信协议ExternalCommandWorker.java 在任意 Trogdor Agent 节点上启动外部命令并通过外部进程的stdin、stdout、stderr与其通信stdout 用于可操作通信stderr 中看到的内容只记录日志。启动时Worker 会先向外部进程发送一条描述工作负载的消息单行 JSON无换行符{id:task ID string, workload:configured workload JSON object}随后监听外部进程发回的 JSON 消息同样为单行。该 JSON 可以包含以下字段status若包含该字段Worker 的状态会被设置为给定值error若包含该字段Worker 的错误会被设置为给定值一旦发生错误外部进程将被终止log若包含该字段将以该文本输出一条日志。一个完整示例{log: Finished successfully., status: {p99ProduceLatency: 100ms, messagesSent: 10000}}协议细节源码注释中明确说明stdout 默认有缓冲外部子进程在写出 status JSON 后最好主动 flush以保证状态能被及时看到如果外部进程向 stdout 输出非 JSON 行Worker 会将其记录为日志如果外部进程退出Worker 结束若退出码非零视为错误需要停止进程时Worker 先发送SIGTERM若在关闭宽限期shutdownGracePeriodMs默认 5000ms见ExternalCommandWorker.DEFAULT_SHUTDOWN_GRACE_PERIOD_MS内未生效则升级为SIGKILL。Exec 模式免 Coordinator 的快速单机测试有时你只想在单个节点上快速跑一个测试不想起 Coordinator。这时可以使用Exec 模式exec mode它允许你在没有 Coordinator 的情况下运行一个单独的 Trogdor Agent。使用 Exec 模式时必须传入一份任务 Spec——Agent 会直接尝试启动该任务./bin/trogdor.sh agent -n node0 -c ./config/trogdor.conf --exec ./tests/spec/simple_produce_bench.jsonExec 模式下任务一旦完成Agent 进程就会退出。从源码看Agent.javaExec 模式内部固定使用 worker IDEXEC_WORKER_ID 1、任务 IDEXEC_TASK_ID task0。这一模式非常适合 CI 流水线中的快速冒烟测试一条命令跑完即退出无需额外清理 Coordinator 进程。深入理解从 Spec 到 Worker 的执行链路把上面的内容串起来一次 Trogdor 任务的完整执行链路是用户通过trogdor.sh client createTask将TaskSpecJSON提交给 Coordinator 的 REST 接口CoordinatorRestResourceCoordinator 的 TaskManager.java 维护任务状态机任务处于PENDING直至startMs到达到达启动时间后Coordinator 调用 Spec 的newController(id)得到TaskController由其决定在哪些节点上创建 worker例如 ProduceBenchSpec 固定返回producerNode单节点Coordinator 向对应 Agent 发送CreateWorkerRequestAgent 的 WorkerManager.java 调用 Spec 的newTaskWorker(id)实例化TaskWorker并启动Worker 在 Agent 内执行具体逻辑生产、消费、故障注入、外部命令等通过WorkerStatusTracker上报状态任务达到durationMs或maxMessages等终止条件后进入STOPPING最终DONE用户通过showTask配合--show-status读取最终状态与延迟统计。整个链路中TaskSpec的newController/newTaskWorker两个抽象方法是Spec → 执行的分叉点见 TaskSpec.java也是扩展自定义任务时唯一需要实现的入口。总结与最佳实践围绕本文内容整理几条实战建议最小环境验证单节点ZooKeeper Broker Agent Coordinator即可完成全流程用 tests/spec/simple_produce_bench.json 跑通 createTask → showTask → showTask --show-status 三步延迟指标解读ProduceBench / ConsumeBench 的输出中p50/p95/p99LatencyMs是评估集群性能的核心transactionsCommitted仅在配置了事务生成器时非零故障演练场景用NetworkPartitionFault验证多 AZ / 多节点场景下的分区容错AutoMQ 强调 Multi-AZ 可用性这类演练尤其有价值用ProcessStopFault模拟节点卡死快速冒烟CI 中优先使用 Exec 模式--exec任务结束 Agent 自动退出无残留进程自定义扩展实现自定义工作负载时继承TaskSpec并实现newController/newTaskWorker两个方法即可接入框架更轻量的方式是借助ExternalCommandSpec直接包装任意脚本或程序通过 stdout JSON 协议上报 status/log/error。Trogdor 将压测、延迟测量、故障注入统一抽象为任务并以Coordinator 编排 Agent 执行 REST/JSON 通信的简洁架构落地是 Kafka 生态中做基准测试与混沌演练的可靠工具。结合本仓库 trogdor/ 的完整源码与测试如 CoordinatorTest.java、ExternalCommandWorkerTest.java你可以进一步深入每个 Spec 的字段语义与 Worker 的线程模型将 Trogdor 无缝融入 AutoMQ 集群的日常质量保障体系。【免费下载链接】automqDiskless Kafka® on S3. 10x Cost-Effective. No Cross-AZ Traffic Cost. Autoscale in seconds. Single-digit ms latency. Multi-AZ Availability.项目地址: https://gitcode.com/GitHub_Trending/au/automq创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表