利用Docker Compose编排实时计算系统需构建低延迟、可协同、状态可控的数据流闭环,涵盖数据接入(Kafka/MQTT)、流处理(Flink等)、状态存储(Redis/RocksDB)和结果展示(Web API/Dashboard),并强调链路职责明确、启动时序保障、资源与日志约束、端到端验证。
要利用 Docker Compose 编排具备实时计算能力的应用系统,关键不是堆砌容器,而是构建一个低延迟、可协同、状态可控的数据流闭环。它通常包含数据接入(如 Kafka 或 MQTT)、流处理引擎(如 Flink、Spark Streaming 或轻量级替代如 Apache Pulsar Functions)、状态存储(Redis 或 RocksDB)、以及结果展示服务(如 Web API 或 Dashboard)。下面从四个实操重点展开:
避免把“实时”简单等同于“快启动”。真正的实时计算系统需明确定义:谁生产数据、谁消费并处理、状态存哪、结果怎么暴露。例如一个传感器告警系统,典型链路是:
confluentinc/cp-kafka 接收设备上报的 MQTT/HTTP 数据,并转为 Kafka Topicflink:1.18-java17 运行自定义 Flink Job(JAR 打包进镜像),做窗口聚合与规则匹配redis:7-alpine 存储滑动窗口计数、最近告警时间戳等轻量状态python:3.11-slim + FastAPI 提供 /alerts 接口,从 Redis 拉取最新结果实时系统对启动顺序和健康就绪非常敏感。不能只靠 depends_on 简单依赖,必须结合 healthcheck 和启动等待逻辑:
healthcheck,检测 kafka-broker 是否响应 describe-cluster
restart: on-failure,防止因 Kafka 尚未就绪导致启动失败退出command 包裹启动脚本,加入 wait-for-it.sh kafka:9092 --timeout=60 -- 等待 Kafka 可达networks 隔离通信实时计算容易因 GC、背压或 OOM 导致延迟飙升,Compose 层需做基础约束:
mem_limit: 2g 和 cpus: '1.5',防止抢占宿主机资源logging.driver: "json-file" 并配置 max-size: "10m",避免日志撑爆磁盘volumes 显式指定 Flink 的 /opt/flink/log 和 checkpoint 目录(如绑定到 ./flink-checkpoints:/opt/flink/checkpoints)redis.conf 自定义配置卷,开启 maxmemory-policy allkeys-lru 防止内存溢出启动后不能只看容器状态为 Up,要验证数据真正流过全链路:
docker-compose exec kafka kafka-console-producer.sh ... 手动发一条测试消息docker-compose logs -f processor 观察 Flink 日志是否打印 Processing record...
docker-compose exec redis redis-cli KEYS "*alert*" 检查 Redis 是否写入预期 keycurl http://localhost:8000/alerts 确认 API 返回最新告警列表docker-compose ps 中容器的 Status 列(是否反复 restarting?是否 Exited with code 137?)