工作流并行度调优
工作流处理机制
工作流通过 Kafka 异步处理任务。事件触发后,系统将工作流消息写入对应的 Kafka Topic,工作流消费服务再读取消息并执行流程。执行期间,流程还会调用工作表等业务服务,并读写数据库。
任务集中触发时,Kafka 会暂存消息,减轻业务服务和数据库的瞬时压力。如果消息积压量持续下降,且任务能在业务允许的时间内处理完成,无须调整并行度。
工作流消费服务使用 md-workflow-consumer 消费组读取以下 Topic:
| Topic | 用途 |
|---|---|
WorkFlow | 主流程执行 |
WorkFlow-Process | 子流程执行 |
WorkFlow-Router | 工作流慢队列 |
WorkFlow-Batch | 批量工作流执行 |
WorkFlow-Button | 按钮触发的工作流执行 |
WorkSheet | 验证行记录是否触发工作流 |
WorkSheet-Batch | 批量验证行记录是否触发工作流 |
WorkSheet-Router | 慢队列验证行记录是否触发工作流 |
一个 Kafka Topic 包含一个或多个分区。在同一个消费组内,一个分区同一时间只能分配给一个消费者。工作流可并行消费的分区数取决于 Topic 分区数、每个消费服务实例的线程数和实例数。
分区数和消费线程总数需要相互匹配。Topic 只有 10 个分区时,即使配置 20 个消费线程,也只有 10 个线程能够获得分区。Topic 有 20 个分区但只有 5 个消费线程时,则只能同时消费 5 个分区。
集群模式的消费线程总数等于每个实例的消费线程数乘以实例数,实际参与消费的线程数不超过 Topic 分区数。
增加 Topic 分区数不会重新分配已有消息。原分区中积压的消息仍在原分区排队;扩容后写入的新消息才会按照分区策略进入新增分区。因此,扩容后原分区的 LAG 不会立即减少。
调整方式
增加工作流并行度前,先确认数据库能否承受更高的并发。数据库 CPU 或磁盘 I/O 使用率较高,或者存在大量慢查询时,增加 Kafka 分区、消费线程和服务实例会继续增加数据库负载。
参考 MongoDB 慢查询优化 检查 慢查询、索引、查询条件和服务器资源。处理数据库瓶颈后,再调整工作流并行度。
单机模式
单机模式使用一个工作流消费进程。增加 Kafka Topic 分区数和消费线程数,可以提高并行处理能力。
所有服务通常部署在同一台服务器上。调整前,检查服务器的 CPU、内存及磁盘 I/O。资源使用率接近上限时,先扩充服务器资源。
相关 Topic 的分区数和工作流消费服务的线程数默认均为 10,调优时建议设置为 20 ~ 30。本例按 20 个分区和 20 个消费线程配置。
-
进入存储组件容器。
docker exec -it $(docker ps | grep mingdaoyun-sc | awk '{print $1}') bash -
查看各 Topic 的分区数和消息积压量。
/usr/local/kafka/bin/kafka-consumer-groups.sh --bootstrap-server 127.0.0.1:9092 --describe --group md-workflow-consumer命令输出中,
PARTITION表示分区编号,LAG表示该分区的消息积压量。 -
根据各 Topic 的积压情况增加分区数。本例将分区数设置为 20。
/usr/local/kafka/bin/kafka-topics.sh --alter --bootstrap-server 127.0.0.1:9092 --partitions 20 --topic WorkFlow/usr/local/kafka/bin/kafka-topics.sh --alter --bootstrap-server 127.0.0.1:9092 --partitions 20 --topic WorkFlow-Batch/usr/local/kafka/bin/kafka-topics.sh --alter --bootstrap-server 127.0.0.1:9092 --partitions 20 --topic WorkFlow-Button/usr/local/kafka/bin/kafka-topics.sh --alter --bootstrap-server 127.0.0.1:9092 --partitions 20 --topic WorkFlow-Process/usr/local/kafka/bin/kafka-topics.sh --alter --bootstrap-server 127.0.0.1:9092 --partitions 20 --topic WorkFlow-Router/usr/local/kafka/bin/kafka-topics.sh --alter --bootstrap-server 127.0.0.1:9092 --partitions 20 --topic WorkSheet/usr/local/kafka/bin/kafka-topics.sh --alter --bootstrap-server 127.0.0.1:9092 --partitions 20 --topic WorkSheet-Batch/usr/local/kafka/bin/kafka-topics.sh --alter --bootstrap-server 127.0.0.1:9092 --partitions 20 --topic WorkSheet-RouterKafka Topic 分区数只能增加,不能减少。根据积压情况分次调整。
-
在
docker-compose.yaml中设置工作流消费线程环境变量。ENV_WORKFLOW_CONSUMER_THREADS: '20'ENV_WORKFLOW_ROUTER_CONSUMER_THREADS: '20'ENV_WORKFLOW_CONSUMER_THREADS控制普通工作流队列的消费线程数,ENV_WORKFLOW_ROUTER_CONSUMER_THREADS控制WorkFlow-Router慢队列的消费线程数。消费线程数通常与对应 Topic 的分区数相同。 -
退出存储组件容器,返回管理器安装目录,重启 HAP 服务使配置生效。
bash ./service.sh restartall
集群模式
集群模式将 HAP 微服务运行在 Kubernetes 中,并独立部署 Kafka。增加 Topic 分区数和消费服务副本数,可以提高并行消费能力。
调整前,检查集群节点的 CPU、内存、磁盘 I/O,以及 Kubernetes Pod 的 CPU 和内存使用率。资源使用率接近上限时,先处理资源瓶颈,再增加并行度。
每个 workflowconsumer 和 workflowrouterconsumer Pod 默认配置 10 个消费线程。Topic 分区数建议设置为 30 ~ 120。本例将分区数设置为 40,并为两个消费服务各配置 4 个 Pod。
-
登录 Kafka 所在服务器,查看各 Topic 的分区数和消息积压量。
/usr/local/kafka/bin/kafka-consumer-groups.sh --bootstrap-server 127.0.0.1:9092 --describe --group md-workflow-consumer命令输出中,
PARTITION表示分区编号,LAG表示该分区的消息积压量。 -
根据各 Topic 的积压情况增加分区数。本例将分区数设置为 40。
/usr/local/kafka/bin/kafka-topics.sh --alter --bootstrap-server 127.0.0.1:9092 --partitions 40 --topic WorkFlow/usr/local/kafka/bin/kafka-topics.sh --alter --bootstrap-server 127.0.0.1:9092 --partitions 40 --topic WorkFlow-Batch/usr/local/kafka/bin/kafka-topics.sh --alter --bootstrap-server 127.0.0.1:9092 --partitions 40 --topic WorkFlow-Button/usr/local/kafka/bin/kafka-topics.sh --alter --bootstrap-server 127.0.0.1:9092 --partitions 40 --topic WorkFlow-Process/usr/local/kafka/bin/kafka-topics.sh --alter --bootstrap-server 127.0.0.1:9092 --partitions 40 --topic WorkFlow-Router/usr/local/kafka/bin/kafka-topics.sh --alter --bootstrap-server 127.0.0.1:9092 --partitions 40 --topic WorkSheet/usr/local/kafka/bin/kafka-topics.sh --alter --bootstrap-server 127.0.0.1:9092 --partitions 40 --topic WorkSheet-Batch/usr/local/kafka/bin/kafka-topics.sh --alter --bootstrap-server 127.0.0.1:9092 --partitions 40 --topic WorkSheet-RouterKafka Topic 分区数只能增加,不能减少。根据积 压情况分次调整。
-
将工作流消费服务扩容到 4 个 Pod。
kubectl scale deployment workflowconsumer --replicas=4kubectl scale deployment workflowrouterconsumer --replicas=4同步修改
service.yaml,将workflowconsumer和workflowrouterconsumer的副本数设为 4,防止 HAP 服务重启后恢复为原副本数。示例:
apiVersion: apps/v1kind: Deploymentmetadata:name: workflowconsumer # <-- 服务名称namespace: defaultspec:replicas: 4 # <-- 副本数workflowconsumer消费普通工作流队列,workflowrouterconsumer消费WorkFlow-Router慢队列。如果修改每个 Pod 的消费线程数,需要根据 Topic 分区数重新计算副本数。 -
提高工作流消费并发后,通过 Kubernetes 集群监控查看工作流所调用微服务的 Pod CPU 使用率。如果某个服务的 CPU 使用率持续较高,增加该服务的副本数。本例将
worksheetonlyworkflow、worksheetonlyworkflowr和command扩容到 4 个 Pod。kubectl scale deployment worksheetonlyworkflow --replicas=4kubectl scale deployment worksheetonlyworkflowr --replicas=4kubectl scale deployment command --replicas=4同步修改
service.yaml,将上述服务的副本数设为 4,防止 HAP 服务重启后恢复为原副本数。