跳到主要内容

工作流并行度调优

工作流处理机制

工作流通过 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 个消费线程配置。

  1. 进入存储组件容器。

    docker exec -it $(docker ps | grep mingdaoyun-sc | awk '{print $1}') bash
  2. 查看各 Topic 的分区数和消息积压量。

    /usr/local/kafka/bin/kafka-consumer-groups.sh --bootstrap-server 127.0.0.1:9092 --describe --group md-workflow-consumer

    命令输出中,PARTITION 表示分区编号,LAG 表示该分区的消息积压量。

  3. 根据各 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-Router

    Kafka Topic 分区数只能增加,不能减少。根据积压情况分次调整。

  4. 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 的分区数相同。

  5. 退出存储组件容器,返回管理器安装目录,重启 HAP 服务使配置生效。

    bash ./service.sh restartall

集群模式

集群模式将 HAP 微服务运行在 Kubernetes 中,并独立部署 Kafka。增加 Topic 分区数和消费服务副本数,可以提高并行消费能力。

调整前,检查集群节点的 CPU、内存、磁盘 I/O,以及 Kubernetes Pod 的 CPU 和内存使用率。资源使用率接近上限时,先处理资源瓶颈,再增加并行度。

每个 workflowconsumerworkflowrouterconsumer Pod 默认配置 10 个消费线程。Topic 分区数建议设置为 30 ~ 120。本例将分区数设置为 40,并为两个消费服务各配置 4 个 Pod。

  1. 登录 Kafka 所在服务器,查看各 Topic 的分区数和消息积压量。

    /usr/local/kafka/bin/kafka-consumer-groups.sh --bootstrap-server 127.0.0.1:9092 --describe --group md-workflow-consumer

    命令输出中,PARTITION 表示分区编号,LAG 表示该分区的消息积压量。

  2. 根据各 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-Router

    Kafka Topic 分区数只能增加,不能减少。根据积压情况分次调整。

  3. 将工作流消费服务扩容到 4 个 Pod。

    kubectl scale deployment workflowconsumer --replicas=4
    kubectl scale deployment workflowrouterconsumer --replicas=4

    同步修改 service.yaml,将 workflowconsumerworkflowrouterconsumer 的副本数设为 4,防止 HAP 服务重启后恢复为原副本数。

    示例:

    apiVersion: apps/v1
    kind: Deployment
    metadata:
    name: workflowconsumer # <-- 服务名称
    namespace: default
    spec:
    replicas: 4 # <-- 副本数

    workflowconsumer 消费普通工作流队列,workflowrouterconsumer 消费 WorkFlow-Router 慢队列。如果修改每个 Pod 的消费线程数,需要根据 Topic 分区数重新计算副本数。

  4. 提高工作流消费并发后,通过 Kubernetes 集群监控查看工作流所调用微服务的 Pod CPU 使用率。如果某个服务的 CPU 使用率持续较高,增加该服务的副本数。本例将 worksheetonlyworkflowworksheetonlyworkflowrcommand 扩容到 4 个 Pod。

    kubectl scale deployment worksheetonlyworkflow --replicas=4
    kubectl scale deployment worksheetonlyworkflowr --replicas=4
    kubectl scale deployment command --replicas=4

    同步修改 service.yaml,将上述服务的副本数设为 4,防止 HAP 服务重启后恢复为原副本数。