From e2c20178cca7cbab0f1bb18548f228b84408abf9 Mon Sep 17 00:00:00 2001 From: JuliaEdom Date: Mon, 20 Jul 2026 14:59:05 +0300 Subject: [PATCH 1/2] fix(stand): bound container logs, restart policies, and Flink restart budget MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit P2.1: x-logging anchor (json-file, 50m x 3) applied to every service in all seven compose files — the 2026-07-18 soak filled the stand disk with unrotated json-file logs (TM alone 3.8G) and took the pipeline down; this makes the ephemeral daemon.json host fix permanent in the repo. P2.2: restart: unless-stopped on flink-jobmanager/flink-taskmanager in docker-compose.flink.yml (one-shot provision jobs keep default "no") — the same soak lost its TM to a container death and hung in NoResourceAvailable for the rest of the run. P2.3: bounded failure-rate restart strategy in both Flink jobs' build_pipeline (3 failures / 5 min, 10s delay, env-overridable) instead of the default infinite fixed-delay — a persistently broken job now parks in FAILED instead of re-emitting from checkpoint forever (soak r1: 116 restarts, 292k duplicates downstream). Unit suite: 2026 passed locally, incl. new restart-strategy tests. Co-Authored-By: Claude Fable 5 --- docker-compose.cdc.yml | 280 ++++---- docker-compose.chaos.yml | 121 ++-- docker-compose.e2e.yml | 368 ++++++----- docker-compose.flink.yml | 163 +++-- docker-compose.iceberg.yml | 115 ++-- docker-compose.prod.yml | 616 +++++++++--------- docker-compose.yml | 464 ++++++------- .../flink_jobs/session_aggregator.py | 11 + src/processing/flink_jobs/stream_processor.py | 14 + tests/unit/test_session_aggregator.py | 30 + tests/unit/test_stream_processor.py | 42 ++ 11 files changed, 1214 insertions(+), 1010 deletions(-) diff --git a/docker-compose.cdc.yml b/docker-compose.cdc.yml index 6925ed87..24fff242 100644 --- a/docker-compose.cdc.yml +++ b/docker-compose.cdc.yml @@ -1,133 +1,147 @@ -services: - kafka: - environment: - KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092 - - cdc-kafka-init: - image: confluentinc/cp-kafka:7.7.0 - depends_on: - kafka: - condition: service_healthy - entrypoint: ["/bin/bash", "-c"] - command: - - | - set -euo pipefail - kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic connect-agentflow-configs --partitions 1 --replication-factor 1 --config cleanup.policy=compact - kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic connect-agentflow-offsets --partitions 25 --replication-factor 1 --config cleanup.policy=compact - kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic connect-agentflow-status --partitions 5 --replication-factor 1 --config cleanup.policy=compact - kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic cdc.postgres.public.orders_v2 --partitions 3 --replication-factor 1 - kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic cdc.postgres.public.users_enriched --partitions 3 --replication-factor 1 - kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic __debezium-heartbeat.cdc.postgres --partitions 1 --replication-factor 1 - kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic cdc.mysql --partitions 1 --replication-factor 1 - kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic __debezium-heartbeat.cdc.mysql --partitions 1 --replication-factor 1 - kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic cdc.mysql.agentflow_demo.products_current --partitions 3 --replication-factor 1 - kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic cdc.mysql.agentflow_demo.sessions_aggregated --partitions 3 --replication-factor 1 - kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic schemahistory.cdc.mysql.agentflow_demo --partitions 1 --replication-factor 1 --config cleanup.policy=delete --config retention.ms=-1 - kafka-configs --bootstrap-server kafka:9092 --alter --entity-type topics --entity-name schemahistory.cdc.mysql.agentflow_demo --add-config cleanup.policy=delete,retention.ms=-1 - - postgres-source: - image: postgres:16 - command: - - "postgres" - - "-c" - - "wal_level=logical" - - "-c" - - "max_replication_slots=4" - - "-c" - - "max_wal_senders=4" - environment: - POSTGRES_DB: agentflow_demo - POSTGRES_USER: cdc_reader - POSTGRES_PASSWORD: agentflow - ports: - - "5433:5432" - volumes: - - ./docker/postgres-source/init.sql:/docker-entrypoint-initdb.d/001-agentflow-cdc.sql:ro - - postgres-cdc-data:/var/lib/postgresql/data - healthcheck: - test: ["CMD-SHELL", "pg_isready -U cdc_reader -d agentflow_demo"] - interval: 10s - timeout: 5s - retries: 5 - - mysql-source: - image: mysql:8.4 - command: - - "--server-id=223344" - - "--log-bin=mysql-bin" - - "--binlog-format=ROW" - - "--binlog-row-image=FULL" - - "--binlog-expire-logs-seconds=864000" - environment: - MYSQL_DATABASE: agentflow_demo - MYSQL_USER: cdc_reader - MYSQL_PASSWORD: agentflow - MYSQL_ROOT_PASSWORD: agentflow-root - ports: - - "3307:3306" - volumes: - - ./docker/mysql-source/init.sql:/docker-entrypoint-initdb.d/001-agentflow-cdc.sql:ro - - mysql-cdc-data:/var/lib/mysql - healthcheck: - test: ["CMD-SHELL", "mysqladmin ping -h localhost -ucdc_reader -pagentflow --silent"] - interval: 10s - timeout: 5s - retries: 10 - - kafka-connect: - image: agentflow-kafka-connect:local - build: - context: . - dockerfile: docker/kafka-connect/Dockerfile - depends_on: - cdc-kafka-init: - condition: service_completed_successfully - postgres-source: - condition: service_healthy - mysql-source: - condition: service_healthy - ports: - - "8083:8083" - environment: - CONNECT_BOOTSTRAP_SERVERS: kafka:9092 - CONNECT_REST_ADVERTISED_HOST_NAME: kafka-connect - CONNECT_REST_PORT: 8083 - CONNECT_GROUP_ID: agentflow-connect - CONNECT_CONFIG_STORAGE_TOPIC: connect-agentflow-configs - CONNECT_OFFSET_STORAGE_TOPIC: connect-agentflow-offsets - CONNECT_STATUS_STORAGE_TOPIC: connect-agentflow-status - CONNECT_CONFIG_STORAGE_REPLICATION_FACTOR: 1 - CONNECT_OFFSET_STORAGE_REPLICATION_FACTOR: 1 - CONNECT_STATUS_STORAGE_REPLICATION_FACTOR: 1 - CONNECT_KEY_CONVERTER: org.apache.kafka.connect.json.JsonConverter - CONNECT_VALUE_CONVERTER: org.apache.kafka.connect.json.JsonConverter - CONNECT_KEY_CONVERTER_SCHEMAS_ENABLE: "false" - CONNECT_VALUE_CONVERTER_SCHEMAS_ENABLE: "false" - CONNECT_PLUGIN_PATH: /usr/share/java,/usr/share/confluent-hub-components,/usr/share/java/debezium - CONNECT_CONFIG_PROVIDERS: file - CONNECT_CONFIG_PROVIDERS_FILE_CLASS: org.apache.kafka.common.config.provider.FileConfigProvider - CONNECT_OFFSET_FLUSH_INTERVAL_MS: 10000 - KAFKA_BOOTSTRAP_SERVERS: kafka:9092 - volumes: - - ./docker/kafka-connect/secrets:/opt/connect/secrets:ro - healthcheck: - test: ["CMD-SHELL", "curl -fsS http://localhost:8083/connectors >/dev/null"] - interval: 10s - timeout: 5s - retries: 12 - - cdc-register-connectors: - image: agentflow-kafka-connect:local - depends_on: - kafka-connect: - condition: service_healthy - environment: - CONNECT_URL: http://kafka-connect:8083 - KAFKA_BOOTSTRAP_SERVERS: kafka:9092 - volumes: - - ./scripts/register_cdc_connectors.sh:/usr/local/bin/register_cdc_connectors.sh:ro - entrypoint: ["/bin/bash", "/usr/local/bin/register_cdc_connectors.sh"] - -volumes: - postgres-cdc-data: - mysql-cdc-data: +# Bounded container logs: the 2026-07-18 soak filled the stand's disk with +# unrotated json-file logs (TM alone grew to 3.8G) and took the pipeline down. +x-logging: &default-logging + driver: json-file + options: + max-size: "50m" + max-file: "3" + +services: + kafka: + logging: *default-logging + environment: + KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092 + + cdc-kafka-init: + logging: *default-logging + image: confluentinc/cp-kafka:7.7.0 + depends_on: + kafka: + condition: service_healthy + entrypoint: ["/bin/bash", "-c"] + command: + - | + set -euo pipefail + kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic connect-agentflow-configs --partitions 1 --replication-factor 1 --config cleanup.policy=compact + kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic connect-agentflow-offsets --partitions 25 --replication-factor 1 --config cleanup.policy=compact + kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic connect-agentflow-status --partitions 5 --replication-factor 1 --config cleanup.policy=compact + kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic cdc.postgres.public.orders_v2 --partitions 3 --replication-factor 1 + kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic cdc.postgres.public.users_enriched --partitions 3 --replication-factor 1 + kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic __debezium-heartbeat.cdc.postgres --partitions 1 --replication-factor 1 + kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic cdc.mysql --partitions 1 --replication-factor 1 + kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic __debezium-heartbeat.cdc.mysql --partitions 1 --replication-factor 1 + kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic cdc.mysql.agentflow_demo.products_current --partitions 3 --replication-factor 1 + kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic cdc.mysql.agentflow_demo.sessions_aggregated --partitions 3 --replication-factor 1 + kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic schemahistory.cdc.mysql.agentflow_demo --partitions 1 --replication-factor 1 --config cleanup.policy=delete --config retention.ms=-1 + kafka-configs --bootstrap-server kafka:9092 --alter --entity-type topics --entity-name schemahistory.cdc.mysql.agentflow_demo --add-config cleanup.policy=delete,retention.ms=-1 + + postgres-source: + logging: *default-logging + image: postgres:16 + command: + - "postgres" + - "-c" + - "wal_level=logical" + - "-c" + - "max_replication_slots=4" + - "-c" + - "max_wal_senders=4" + environment: + POSTGRES_DB: agentflow_demo + POSTGRES_USER: cdc_reader + POSTGRES_PASSWORD: agentflow + ports: + - "5433:5432" + volumes: + - ./docker/postgres-source/init.sql:/docker-entrypoint-initdb.d/001-agentflow-cdc.sql:ro + - postgres-cdc-data:/var/lib/postgresql/data + healthcheck: + test: ["CMD-SHELL", "pg_isready -U cdc_reader -d agentflow_demo"] + interval: 10s + timeout: 5s + retries: 5 + + mysql-source: + logging: *default-logging + image: mysql:8.4 + command: + - "--server-id=223344" + - "--log-bin=mysql-bin" + - "--binlog-format=ROW" + - "--binlog-row-image=FULL" + - "--binlog-expire-logs-seconds=864000" + environment: + MYSQL_DATABASE: agentflow_demo + MYSQL_USER: cdc_reader + MYSQL_PASSWORD: agentflow + MYSQL_ROOT_PASSWORD: agentflow-root + ports: + - "3307:3306" + volumes: + - ./docker/mysql-source/init.sql:/docker-entrypoint-initdb.d/001-agentflow-cdc.sql:ro + - mysql-cdc-data:/var/lib/mysql + healthcheck: + test: ["CMD-SHELL", "mysqladmin ping -h localhost -ucdc_reader -pagentflow --silent"] + interval: 10s + timeout: 5s + retries: 10 + + kafka-connect: + logging: *default-logging + image: agentflow-kafka-connect:local + build: + context: . + dockerfile: docker/kafka-connect/Dockerfile + depends_on: + cdc-kafka-init: + condition: service_completed_successfully + postgres-source: + condition: service_healthy + mysql-source: + condition: service_healthy + ports: + - "8083:8083" + environment: + CONNECT_BOOTSTRAP_SERVERS: kafka:9092 + CONNECT_REST_ADVERTISED_HOST_NAME: kafka-connect + CONNECT_REST_PORT: 8083 + CONNECT_GROUP_ID: agentflow-connect + CONNECT_CONFIG_STORAGE_TOPIC: connect-agentflow-configs + CONNECT_OFFSET_STORAGE_TOPIC: connect-agentflow-offsets + CONNECT_STATUS_STORAGE_TOPIC: connect-agentflow-status + CONNECT_CONFIG_STORAGE_REPLICATION_FACTOR: 1 + CONNECT_OFFSET_STORAGE_REPLICATION_FACTOR: 1 + CONNECT_STATUS_STORAGE_REPLICATION_FACTOR: 1 + CONNECT_KEY_CONVERTER: org.apache.kafka.connect.json.JsonConverter + CONNECT_VALUE_CONVERTER: org.apache.kafka.connect.json.JsonConverter + CONNECT_KEY_CONVERTER_SCHEMAS_ENABLE: "false" + CONNECT_VALUE_CONVERTER_SCHEMAS_ENABLE: "false" + CONNECT_PLUGIN_PATH: /usr/share/java,/usr/share/confluent-hub-components,/usr/share/java/debezium + CONNECT_CONFIG_PROVIDERS: file + CONNECT_CONFIG_PROVIDERS_FILE_CLASS: org.apache.kafka.common.config.provider.FileConfigProvider + CONNECT_OFFSET_FLUSH_INTERVAL_MS: 10000 + KAFKA_BOOTSTRAP_SERVERS: kafka:9092 + volumes: + - ./docker/kafka-connect/secrets:/opt/connect/secrets:ro + healthcheck: + test: ["CMD-SHELL", "curl -fsS http://localhost:8083/connectors >/dev/null"] + interval: 10s + timeout: 5s + retries: 12 + + cdc-register-connectors: + logging: *default-logging + image: agentflow-kafka-connect:local + depends_on: + kafka-connect: + condition: service_healthy + environment: + CONNECT_URL: http://kafka-connect:8083 + KAFKA_BOOTSTRAP_SERVERS: kafka:9092 + volumes: + - ./scripts/register_cdc_connectors.sh:/usr/local/bin/register_cdc_connectors.sh:ro + entrypoint: ["/bin/bash", "/usr/local/bin/register_cdc_connectors.sh"] + +volumes: + postgres-cdc-data: + mysql-cdc-data: diff --git a/docker-compose.chaos.yml b/docker-compose.chaos.yml index 93c71bf8..9107da70 100644 --- a/docker-compose.chaos.yml +++ b/docker-compose.chaos.yml @@ -1,55 +1,66 @@ -services: - kafka: - image: confluentinc/cp-kafka:7.7.0 - hostname: kafka - environment: - CLUSTER_ID: AgentFlowChaosCluster01 - KAFKA_NODE_ID: 1 - KAFKA_PROCESS_ROLES: broker,controller - KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093 - KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://127.0.0.1:${AGENTFLOW_CHAOS_KAFKA_PORT:-19092} - KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT - KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER - KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT - KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:9093 - KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 - KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 - KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 - KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true" - healthcheck: - test: ["CMD-SHELL", "kafka-broker-api-versions --bootstrap-server localhost:9092 >/dev/null 2>&1"] - interval: 5s - timeout: 5s - retries: 24 - start_period: 20s - - redis: - image: redis:7.4-alpine - command: ["redis-server", "--save", "", "--appendonly", "no"] - healthcheck: - test: ["CMD", "redis-cli", "ping"] - interval: 5s - timeout: 3s - retries: 24 - start_period: 5s - - toxiproxy: - image: ghcr.io/shopify/toxiproxy:2.9.0 - depends_on: - kafka: - condition: service_healthy - redis: - condition: service_healthy - command: ["-host", "0.0.0.0", "-config", "/config/toxiproxy.json"] - volumes: - - ./config/toxiproxy.json:/config/toxiproxy.json:ro - ports: - - "${AGENTFLOW_CHAOS_TOXIPROXY_PORT:-8474}:8474" - - "${AGENTFLOW_CHAOS_KAFKA_PORT:-19092}:19092" - - "${AGENTFLOW_CHAOS_REDIS_PORT:-16380}:16380" - healthcheck: - test: ["CMD", "/toxiproxy-cli", "list"] - interval: 5s - timeout: 5s - retries: 24 - start_period: 5s +# Bounded container logs: the 2026-07-18 soak filled the stand's disk with +# unrotated json-file logs (TM alone grew to 3.8G) and took the pipeline down. +x-logging: &default-logging + driver: json-file + options: + max-size: "50m" + max-file: "3" + +services: + kafka: + logging: *default-logging + image: confluentinc/cp-kafka:7.7.0 + hostname: kafka + environment: + CLUSTER_ID: AgentFlowChaosCluster01 + KAFKA_NODE_ID: 1 + KAFKA_PROCESS_ROLES: broker,controller + KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093 + KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://127.0.0.1:${AGENTFLOW_CHAOS_KAFKA_PORT:-19092} + KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT + KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER + KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT + KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:9093 + KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 + KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 + KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 + KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true" + healthcheck: + test: ["CMD-SHELL", "kafka-broker-api-versions --bootstrap-server localhost:9092 >/dev/null 2>&1"] + interval: 5s + timeout: 5s + retries: 24 + start_period: 20s + + redis: + logging: *default-logging + image: redis:7.4-alpine + command: ["redis-server", "--save", "", "--appendonly", "no"] + healthcheck: + test: ["CMD", "redis-cli", "ping"] + interval: 5s + timeout: 3s + retries: 24 + start_period: 5s + + toxiproxy: + logging: *default-logging + image: ghcr.io/shopify/toxiproxy:2.9.0 + depends_on: + kafka: + condition: service_healthy + redis: + condition: service_healthy + command: ["-host", "0.0.0.0", "-config", "/config/toxiproxy.json"] + volumes: + - ./config/toxiproxy.json:/config/toxiproxy.json:ro + ports: + - "${AGENTFLOW_CHAOS_TOXIPROXY_PORT:-8474}:8474" + - "${AGENTFLOW_CHAOS_KAFKA_PORT:-19092}:19092" + - "${AGENTFLOW_CHAOS_REDIS_PORT:-16380}:16380" + healthcheck: + test: ["CMD", "/toxiproxy-cli", "list"] + interval: 5s + timeout: 5s + retries: 24 + start_period: 5s diff --git a/docker-compose.e2e.yml b/docker-compose.e2e.yml index d406c279..0120e640 100644 --- a/docker-compose.e2e.yml +++ b/docker-compose.e2e.yml @@ -1,177 +1,191 @@ -services: - redis: - image: redis:7.4-alpine - command: ["redis-server", "--save", "", "--appendonly", "no"] - healthcheck: - test: ["CMD", "redis-cli", "ping"] - interval: 5s - timeout: 5s - retries: 20 - deploy: - resources: - limits: - cpus: "0.50" - memory: 256M - - kafka: - image: confluentinc/cp-kafka:7.7.0 - hostname: kafka - ports: - - "19092:19092" - environment: - CLUSTER_ID: "MkU3OEVBNTcwNTJENDM2Qk" - KAFKA_NODE_ID: 1 - KAFKA_PROCESS_ROLES: broker,controller - KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:29093,EXTERNAL://0.0.0.0:19092 - KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,EXTERNAL://localhost:19092 - KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT,EXTERNAL:PLAINTEXT - KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT - KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER - KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:29093 - KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true" - KAFKA_DEFAULT_REPLICATION_FACTOR: 1 - KAFKA_MIN_INSYNC_REPLICAS: 1 - KAFKA_NUM_PARTITIONS: 3 - KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 - KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 - KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 - KAFKA_CONFLUENT_LICENSE_TOPIC_REPLICATION_FACTOR: 1 - KAFKA_LOG_DIRS: /var/lib/kafka/data - KAFKA_HEAP_OPTS: "-Xms256m -Xmx512m" - healthcheck: - test: ["CMD-SHELL", "kafka-topics --bootstrap-server localhost:9092 --list >/dev/null 2>&1"] - interval: 10s - timeout: 10s - retries: 12 - start_period: 15s - deploy: - resources: - limits: - cpus: "1.00" - memory: 1G - - clickhouse: - image: clickhouse/clickhouse-server:24.8 - environment: - CLICKHOUSE_DB: agentflow - CLICKHOUSE_USER: default - CLICKHOUSE_PASSWORD: "" - CLICKHOUSE_DEFAULT_ACCESS_MANAGEMENT: "1" - healthcheck: - test: ["CMD-SHELL", "wget -qO- http://localhost:8123/ping | grep -q 'Ok.'"] - interval: 10s - timeout: 10s - retries: 12 - start_period: 10s - deploy: - resources: - limits: - cpus: "1.00" - memory: 1G - - postgres: - image: postgres:16-alpine - environment: - POSTGRES_DB: agentflow - POSTGRES_USER: agentflow - POSTGRES_PASSWORD: agentflow - healthcheck: - test: ["CMD-SHELL", "pg_isready -U agentflow -d agentflow"] - interval: 5s - timeout: 5s - retries: 20 - deploy: - resources: - limits: - cpus: "0.50" - memory: 512M - - serving-init: - # The API stopped provisioning and seeding the serving store on boot - # (audit P0-2), so the stack lays the schema and the documented demo - # entities down itself before the API starts — the same one-shot job - # docker-compose.prod.yml runs, against the in-stack ClickHouse the E2E - # smoke suite queries. - build: - context: . - dockerfile: Dockerfile.api - depends_on: - clickhouse: - condition: service_healthy - environment: - AGENTFLOW_SERVING_CONFIG: /app/config/serving.yaml - DUCKDB_PATH: /app/data/agentflow.duckdb - CLICKHOUSE_HOST: clickhouse - CLICKHOUSE_PORT: 8123 - CLICKHOUSE_DATABASE: agentflow - CLICKHOUSE_USER: default - CLICKHOUSE_PASSWORD: "" - command: ["python", "-m", "src.serving.provision", "--schema", "--seed"] - restart: "no" - volumes: - - agentflow-api-data:/app/data - - ./config:/app/config:ro - - agentflow-api: - build: - context: . - dockerfile: Dockerfile.api - depends_on: - redis: - condition: service_healthy - kafka: - condition: service_healthy - postgres: - condition: service_healthy - clickhouse: - condition: service_healthy - serving-init: - condition: service_completed_successfully - ports: - - "8000:8000" - environment: - AGENTFLOW_SERVING_CONFIG: /app/config/serving.yaml - AGENTFLOW_API_KEYS: ${AGENTFLOW_API_KEYS:-af-prod-agent-support-abc123:Support Agent,af-prod-agent-ops-def456:Ops Agent,af-prod-agent-rate-e2e000:Rate Limit Agent} - AGENTFLOW_RATE_LIMIT_RPM: ${AGENTFLOW_RATE_LIMIT_RPM:-120} - AGENTFLOW_USAGE_DB_PATH: /app/data/agentflow_api_usage.duckdb - AGENTFLOW_WEBHOOKS_FILE: /app/data/e2e-webhooks.yaml - # The e2e webhook callback is delivered to the host gateway, which the - # SSRF egress guard rejects as a private address. Trust exactly that - # configured callback host (default host.docker.internal) — the guard - # stays fully active for every other target. - AGENTFLOW_EGRESS_ALLOWED_HOSTS: ${AGENTFLOW_E2E_CALLBACK_HOST:-host.docker.internal} - DUCKDB_PATH: /app/data/agentflow.duckdb - # E2E runs the shipped serving profile: config/serving.yaml selects - # ClickHouse, and these point the backend at the in-stack server. - CLICKHOUSE_HOST: clickhouse - CLICKHOUSE_PORT: 8123 - CLICKHOUSE_DATABASE: agentflow - CLICKHOUSE_USER: default - CLICKHOUSE_PASSWORD: "" - REDIS_URL: redis://redis:6379/0 - KAFKA_BOOTSTRAP_SERVERS: kafka:9092 - OTEL_SDK_DISABLED: "true" - extra_hosts: - - "host.docker.internal:host-gateway" - volumes: - - agentflow-api-data:/app/data - - ./config:/app/config:ro - - ./contracts:/app/contracts:ro - healthcheck: - test: - - CMD - - python - - -c - - import urllib.request; urllib.request.urlopen("http://localhost:8000/v1/health", timeout=5) - interval: 10s - timeout: 10s - retries: 12 - start_period: 15s - deploy: - resources: - limits: - cpus: "1.50" - memory: 2G - -volumes: - agentflow-api-data: +# Bounded container logs: the 2026-07-18 soak filled the stand's disk with +# unrotated json-file logs (TM alone grew to 3.8G) and took the pipeline down. +x-logging: &default-logging + driver: json-file + options: + max-size: "50m" + max-file: "3" + +services: + redis: + logging: *default-logging + image: redis:7.4-alpine + command: ["redis-server", "--save", "", "--appendonly", "no"] + healthcheck: + test: ["CMD", "redis-cli", "ping"] + interval: 5s + timeout: 5s + retries: 20 + deploy: + resources: + limits: + cpus: "0.50" + memory: 256M + + kafka: + logging: *default-logging + image: confluentinc/cp-kafka:7.7.0 + hostname: kafka + ports: + - "19092:19092" + environment: + CLUSTER_ID: "MkU3OEVBNTcwNTJENDM2Qk" + KAFKA_NODE_ID: 1 + KAFKA_PROCESS_ROLES: broker,controller + KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:29093,EXTERNAL://0.0.0.0:19092 + KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,EXTERNAL://localhost:19092 + KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT,EXTERNAL:PLAINTEXT + KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT + KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER + KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:29093 + KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true" + KAFKA_DEFAULT_REPLICATION_FACTOR: 1 + KAFKA_MIN_INSYNC_REPLICAS: 1 + KAFKA_NUM_PARTITIONS: 3 + KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 + KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 + KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 + KAFKA_CONFLUENT_LICENSE_TOPIC_REPLICATION_FACTOR: 1 + KAFKA_LOG_DIRS: /var/lib/kafka/data + KAFKA_HEAP_OPTS: "-Xms256m -Xmx512m" + healthcheck: + test: ["CMD-SHELL", "kafka-topics --bootstrap-server localhost:9092 --list >/dev/null 2>&1"] + interval: 10s + timeout: 10s + retries: 12 + start_period: 15s + deploy: + resources: + limits: + cpus: "1.00" + memory: 1G + + clickhouse: + logging: *default-logging + image: clickhouse/clickhouse-server:24.8 + environment: + CLICKHOUSE_DB: agentflow + CLICKHOUSE_USER: default + CLICKHOUSE_PASSWORD: "" + CLICKHOUSE_DEFAULT_ACCESS_MANAGEMENT: "1" + healthcheck: + test: ["CMD-SHELL", "wget -qO- http://localhost:8123/ping | grep -q 'Ok.'"] + interval: 10s + timeout: 10s + retries: 12 + start_period: 10s + deploy: + resources: + limits: + cpus: "1.00" + memory: 1G + + postgres: + logging: *default-logging + image: postgres:16-alpine + environment: + POSTGRES_DB: agentflow + POSTGRES_USER: agentflow + POSTGRES_PASSWORD: agentflow + healthcheck: + test: ["CMD-SHELL", "pg_isready -U agentflow -d agentflow"] + interval: 5s + timeout: 5s + retries: 20 + deploy: + resources: + limits: + cpus: "0.50" + memory: 512M + + serving-init: + logging: *default-logging + # The API stopped provisioning and seeding the serving store on boot + # (audit P0-2), so the stack lays the schema and the documented demo + # entities down itself before the API starts — the same one-shot job + # docker-compose.prod.yml runs, against the in-stack ClickHouse the E2E + # smoke suite queries. + build: + context: . + dockerfile: Dockerfile.api + depends_on: + clickhouse: + condition: service_healthy + environment: + AGENTFLOW_SERVING_CONFIG: /app/config/serving.yaml + DUCKDB_PATH: /app/data/agentflow.duckdb + CLICKHOUSE_HOST: clickhouse + CLICKHOUSE_PORT: 8123 + CLICKHOUSE_DATABASE: agentflow + CLICKHOUSE_USER: default + CLICKHOUSE_PASSWORD: "" + command: ["python", "-m", "src.serving.provision", "--schema", "--seed"] + restart: "no" + volumes: + - agentflow-api-data:/app/data + - ./config:/app/config:ro + + agentflow-api: + logging: *default-logging + build: + context: . + dockerfile: Dockerfile.api + depends_on: + redis: + condition: service_healthy + kafka: + condition: service_healthy + postgres: + condition: service_healthy + clickhouse: + condition: service_healthy + serving-init: + condition: service_completed_successfully + ports: + - "8000:8000" + environment: + AGENTFLOW_SERVING_CONFIG: /app/config/serving.yaml + AGENTFLOW_API_KEYS: ${AGENTFLOW_API_KEYS:-af-prod-agent-support-abc123:Support Agent,af-prod-agent-ops-def456:Ops Agent,af-prod-agent-rate-e2e000:Rate Limit Agent} + AGENTFLOW_RATE_LIMIT_RPM: ${AGENTFLOW_RATE_LIMIT_RPM:-120} + AGENTFLOW_USAGE_DB_PATH: /app/data/agentflow_api_usage.duckdb + AGENTFLOW_WEBHOOKS_FILE: /app/data/e2e-webhooks.yaml + # The e2e webhook callback is delivered to the host gateway, which the + # SSRF egress guard rejects as a private address. Trust exactly that + # configured callback host (default host.docker.internal) — the guard + # stays fully active for every other target. + AGENTFLOW_EGRESS_ALLOWED_HOSTS: ${AGENTFLOW_E2E_CALLBACK_HOST:-host.docker.internal} + DUCKDB_PATH: /app/data/agentflow.duckdb + # E2E runs the shipped serving profile: config/serving.yaml selects + # ClickHouse, and these point the backend at the in-stack server. + CLICKHOUSE_HOST: clickhouse + CLICKHOUSE_PORT: 8123 + CLICKHOUSE_DATABASE: agentflow + CLICKHOUSE_USER: default + CLICKHOUSE_PASSWORD: "" + REDIS_URL: redis://redis:6379/0 + KAFKA_BOOTSTRAP_SERVERS: kafka:9092 + OTEL_SDK_DISABLED: "true" + extra_hosts: + - "host.docker.internal:host-gateway" + volumes: + - agentflow-api-data:/app/data + - ./config:/app/config:ro + - ./contracts:/app/contracts:ro + healthcheck: + test: + - CMD + - python + - -c + - import urllib.request; urllib.request.urlopen("http://localhost:8000/v1/health", timeout=5) + interval: 10s + timeout: 10s + retries: 12 + start_period: 15s + deploy: + resources: + limits: + cpus: "1.50" + memory: 2G + +volumes: + agentflow-api-data: diff --git a/docker-compose.flink.yml b/docker-compose.flink.yml index 3d9a0729..838ccf89 100644 --- a/docker-compose.flink.yml +++ b/docker-compose.flink.yml @@ -1,72 +1,91 @@ -services: - kafka: - environment: - KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092,HOST://0.0.0.0:19092,CONTROLLER://0.0.0.0:29093 - KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,HOST://localhost:19092 - KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,HOST:PLAINTEXT,CONTROLLER:PLAINTEXT - KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT - ports: - - "9092:9092" - - "19092:19092" - - "9101:9101" - - flink-jobmanager: - build: - context: ./src - dockerfile: processing/flink_jobs/Dockerfile - image: agentflow-flink-local:latest - command: - - /bin/bash - - -lc - - | - printf '%s\n' "$$FLINK_PROPERTIES" >> /opt/flink/conf/config.yaml - printf '%s\n' 'jobmanager.memory.process.size: 1600m' >> /opt/flink/conf/config.yaml - exec /opt/flink/bin/jobmanager.sh start-foreground - environment: - PYTHONPATH: /opt/agentflow - - flink-taskmanager: - build: - context: ./src - dockerfile: processing/flink_jobs/Dockerfile - image: agentflow-flink-local:latest - command: - - /bin/bash - - -lc - - | - printf '%s\n' "$$FLINK_PROPERTIES" >> /opt/flink/conf/config.yaml - printf '%s\n' 'taskmanager.memory.process.size: 1728m' >> /opt/flink/conf/config.yaml - exec /opt/flink/bin/taskmanager.sh start-foreground - environment: - PYTHONPATH: /opt/agentflow - - flink-job-runner: - build: - context: ./src - dockerfile: processing/flink_jobs/Dockerfile - image: agentflow-flink-local:latest - depends_on: - kafka-init: - condition: service_completed_successfully - minio-init: - condition: service_completed_successfully - flink-jobmanager: - condition: service_healthy - flink-taskmanager: - condition: service_started - working_dir: /opt/agentflow - environment: - KAFKA_BOOTSTRAP_SERVERS: kafka:9092 - FLINK_PARALLELISM: "2" - PYTHONPATH: /opt/agentflow - command: - - /bin/bash - - -lc - - | - set -euo pipefail - until curl -fsS http://flink-jobmanager:8081/overview >/dev/null; do - sleep 2 - done - /opt/flink/bin/flink run -d -m flink-jobmanager:8081 -py /opt/agentflow/src/processing/flink_jobs/stream_processor.py - echo "Submitted src/processing/flink_jobs/stream_processor.py to the local Flink cluster." - tail -f /dev/null +# Bounded container logs: the 2026-07-18 soak filled the stand's disk with +# unrotated json-file logs (TM alone grew to 3.8G) and took the pipeline down. +x-logging: &default-logging + driver: json-file + options: + max-size: "50m" + max-file: "3" + +services: + kafka: + logging: *default-logging + environment: + KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092,HOST://0.0.0.0:19092,CONTROLLER://0.0.0.0:29093 + KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,HOST://localhost:19092 + KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,HOST:PLAINTEXT,CONTROLLER:PLAINTEXT + KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT + ports: + - "9092:9092" + - "19092:19092" + - "9101:9101" + + flink-jobmanager: + logging: *default-logging + # The 2026-07-18 soak lost its TaskManager to a one-off container death and + # the job then hung in NoResourceAvailable for the rest of the run; the + # cluster processes must come back on their own. The one-shot job + # containers (kafka-init, minio-init, flink-job-runner) keep the default + # "no" — restarting a completed provision job would resubmit the Flink job. + restart: unless-stopped + build: + context: ./src + dockerfile: processing/flink_jobs/Dockerfile + image: agentflow-flink-local:latest + command: + - /bin/bash + - -lc + - | + printf '%s\n' "$$FLINK_PROPERTIES" >> /opt/flink/conf/config.yaml + printf '%s\n' 'jobmanager.memory.process.size: 1600m' >> /opt/flink/conf/config.yaml + exec /opt/flink/bin/jobmanager.sh start-foreground + environment: + PYTHONPATH: /opt/agentflow + + flink-taskmanager: + logging: *default-logging + restart: unless-stopped + build: + context: ./src + dockerfile: processing/flink_jobs/Dockerfile + image: agentflow-flink-local:latest + command: + - /bin/bash + - -lc + - | + printf '%s\n' "$$FLINK_PROPERTIES" >> /opt/flink/conf/config.yaml + printf '%s\n' 'taskmanager.memory.process.size: 1728m' >> /opt/flink/conf/config.yaml + exec /opt/flink/bin/taskmanager.sh start-foreground + environment: + PYTHONPATH: /opt/agentflow + + flink-job-runner: + logging: *default-logging + build: + context: ./src + dockerfile: processing/flink_jobs/Dockerfile + image: agentflow-flink-local:latest + depends_on: + kafka-init: + condition: service_completed_successfully + minio-init: + condition: service_completed_successfully + flink-jobmanager: + condition: service_healthy + flink-taskmanager: + condition: service_started + working_dir: /opt/agentflow + environment: + KAFKA_BOOTSTRAP_SERVERS: kafka:9092 + FLINK_PARALLELISM: "2" + PYTHONPATH: /opt/agentflow + command: + - /bin/bash + - -lc + - | + set -euo pipefail + until curl -fsS http://flink-jobmanager:8081/overview >/dev/null; do + sleep 2 + done + /opt/flink/bin/flink run -d -m flink-jobmanager:8081 -py /opt/agentflow/src/processing/flink_jobs/stream_processor.py + echo "Submitted src/processing/flink_jobs/stream_processor.py to the local Flink cluster." + tail -f /dev/null diff --git a/docker-compose.iceberg.yml b/docker-compose.iceberg.yml index 9207e6f5..abb21a7e 100644 --- a/docker-compose.iceberg.yml +++ b/docker-compose.iceberg.yml @@ -1,52 +1,63 @@ -# Focused Iceberg REST catalog stack, backed by MinIO (S3) instead of an -# ephemeral local directory. Mirrors the object-store wiring in -# docker-compose.yml (same image tags, bucket, credentials) so the PyIceberg -# sink writes table data/metadata to the real object store the rest of the -# stack already uses. Consumed by tests/integration/test_iceberg_sink.py. -services: - minio: - image: minio/minio:RELEASE.2025-09-07T16-13-09Z - ports: - - "9000:9000" - - "9001:9001" - command: server /data --console-address ":9001" - environment: - MINIO_ROOT_USER: minio - MINIO_ROOT_PASSWORD: minio123 - healthcheck: - test: ["CMD", "curl", "-f", "http://localhost:9000/minio/health/live"] - interval: 5s - timeout: 5s - retries: 5 - - # One-shot bucket provisioner. Retries the alias instead of gating on the - # MinIO healthcheck so the stack comes up regardless of whether the server - # image ships curl. - minio-init: - image: minio/mc:RELEASE.2025-08-13T08-35-41Z - depends_on: - - minio - entrypoint: ["/bin/bash", "-c"] - command: - - | - until mc alias set local http://minio:9000 minio minio123; do - echo "waiting for minio..."; sleep 1; - done - mc mb local/agentflow-lake --ignore-existing - echo "MinIO bucket ready." - - iceberg-rest: - image: tabulario/iceberg-rest:0.6.0 - depends_on: - minio-init: - condition: service_completed_successfully - ports: - - "8181:8181" - environment: - CATALOG_WAREHOUSE: s3://agentflow-lake/warehouse - CATALOG_IO__IMPL: org.apache.iceberg.aws.s3.S3FileIO - CATALOG_S3_ENDPOINT: http://minio:9000 - CATALOG_S3_PATH__STYLE__ACCESS: "true" - AWS_ACCESS_KEY_ID: minio - AWS_SECRET_ACCESS_KEY: minio123 - AWS_REGION: us-east-1 +# Focused Iceberg REST catalog stack, backed by MinIO (S3) instead of an +# ephemeral local directory. Mirrors the object-store wiring in +# docker-compose.yml (same image tags, bucket, credentials) so the PyIceberg +# sink writes table data/metadata to the real object store the rest of the +# stack already uses. Consumed by tests/integration/test_iceberg_sink.py. +# Bounded container logs: the 2026-07-18 soak filled the stand's disk with +# unrotated json-file logs (TM alone grew to 3.8G) and took the pipeline down. +x-logging: &default-logging + driver: json-file + options: + max-size: "50m" + max-file: "3" + +services: + minio: + logging: *default-logging + image: minio/minio:RELEASE.2025-09-07T16-13-09Z + ports: + - "9000:9000" + - "9001:9001" + command: server /data --console-address ":9001" + environment: + MINIO_ROOT_USER: minio + MINIO_ROOT_PASSWORD: minio123 + healthcheck: + test: ["CMD", "curl", "-f", "http://localhost:9000/minio/health/live"] + interval: 5s + timeout: 5s + retries: 5 + + # One-shot bucket provisioner. Retries the alias instead of gating on the + # MinIO healthcheck so the stack comes up regardless of whether the server + # image ships curl. + minio-init: + logging: *default-logging + image: minio/mc:RELEASE.2025-08-13T08-35-41Z + depends_on: + - minio + entrypoint: ["/bin/bash", "-c"] + command: + - | + until mc alias set local http://minio:9000 minio minio123; do + echo "waiting for minio..."; sleep 1; + done + mc mb local/agentflow-lake --ignore-existing + echo "MinIO bucket ready." + + iceberg-rest: + logging: *default-logging + image: tabulario/iceberg-rest:0.6.0 + depends_on: + minio-init: + condition: service_completed_successfully + ports: + - "8181:8181" + environment: + CATALOG_WAREHOUSE: s3://agentflow-lake/warehouse + CATALOG_IO__IMPL: org.apache.iceberg.aws.s3.S3FileIO + CATALOG_S3_ENDPOINT: http://minio:9000 + CATALOG_S3_PATH__STYLE__ACCESS: "true" + AWS_ACCESS_KEY_ID: minio + AWS_SECRET_ACCESS_KEY: minio123 + AWS_REGION: us-east-1 diff --git a/docker-compose.prod.yml b/docker-compose.prod.yml index 7c009023..d479a563 100644 --- a/docker-compose.prod.yml +++ b/docker-compose.prod.yml @@ -1,298 +1,318 @@ -services: - kafka-1: - image: confluentinc/cp-kafka:7.7.0 - hostname: kafka-1 - ports: - - "127.0.0.1:19092:19092" - environment: - CLUSTER_ID: "AgentFlowKafkaProdCluster01" - KAFKA_NODE_ID: 1 - KAFKA_PROCESS_ROLES: broker,controller - KAFKA_LISTENERS: INTERNAL://0.0.0.0:9092,EXTERNAL://0.0.0.0:19092,CONTROLLER://0.0.0.0:29093 - KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka-1:9092,EXTERNAL://localhost:19092 - KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT,CONTROLLER:PLAINTEXT - KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL - KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER - KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka-1:29093,2@kafka-2:29093,3@kafka-3:29093 - KAFKA_LOG_DIRS: /var/lib/kafka/data - KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true" - KAFKA_DEFAULT_REPLICATION_FACTOR: 3 - KAFKA_MIN_INSYNC_REPLICAS: 2 - KAFKA_NUM_PARTITIONS: 6 - KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3 - KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3 - KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2 - KAFKA_CONFLUENT_LICENSE_TOPIC_REPLICATION_FACTOR: 3 - volumes: - - kafka-1-data:/var/lib/kafka/data - healthcheck: - test: ["CMD-SHELL", "kafka-broker-api-versions --bootstrap-server localhost:9092"] - interval: 10s - timeout: 10s - retries: 10 - - kafka-2: - image: confluentinc/cp-kafka:7.7.0 - hostname: kafka-2 - ports: - - "127.0.0.1:29092:19092" - environment: - CLUSTER_ID: "AgentFlowKafkaProdCluster01" - KAFKA_NODE_ID: 2 - KAFKA_PROCESS_ROLES: broker,controller - KAFKA_LISTENERS: INTERNAL://0.0.0.0:9092,EXTERNAL://0.0.0.0:19092,CONTROLLER://0.0.0.0:29093 - KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka-2:9092,EXTERNAL://localhost:29092 - KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT,CONTROLLER:PLAINTEXT - KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL - KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER - KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka-1:29093,2@kafka-2:29093,3@kafka-3:29093 - KAFKA_LOG_DIRS: /var/lib/kafka/data - KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true" - KAFKA_DEFAULT_REPLICATION_FACTOR: 3 - KAFKA_MIN_INSYNC_REPLICAS: 2 - KAFKA_NUM_PARTITIONS: 6 - KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3 - KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3 - KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2 - KAFKA_CONFLUENT_LICENSE_TOPIC_REPLICATION_FACTOR: 3 - volumes: - - kafka-2-data:/var/lib/kafka/data - healthcheck: - test: ["CMD-SHELL", "kafka-broker-api-versions --bootstrap-server localhost:9092"] - interval: 10s - timeout: 10s - retries: 10 - - kafka-3: - image: confluentinc/cp-kafka:7.7.0 - hostname: kafka-3 - ports: - - "127.0.0.1:39092:19092" - environment: - CLUSTER_ID: "AgentFlowKafkaProdCluster01" - KAFKA_NODE_ID: 3 - KAFKA_PROCESS_ROLES: broker,controller - KAFKA_LISTENERS: INTERNAL://0.0.0.0:9092,EXTERNAL://0.0.0.0:19092,CONTROLLER://0.0.0.0:29093 - KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka-3:9092,EXTERNAL://localhost:39092 - KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT,CONTROLLER:PLAINTEXT - KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL - KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER - KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka-1:29093,2@kafka-2:29093,3@kafka-3:29093 - KAFKA_LOG_DIRS: /var/lib/kafka/data - KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true" - KAFKA_DEFAULT_REPLICATION_FACTOR: 3 - KAFKA_MIN_INSYNC_REPLICAS: 2 - KAFKA_NUM_PARTITIONS: 6 - KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3 - KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3 - KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2 - KAFKA_CONFLUENT_LICENSE_TOPIC_REPLICATION_FACTOR: 3 - volumes: - - kafka-3-data:/var/lib/kafka/data - healthcheck: - test: ["CMD-SHELL", "kafka-broker-api-versions --bootstrap-server localhost:9092"] - interval: 10s - timeout: 10s - retries: 10 - - schema-registry: - image: confluentinc/cp-schema-registry:7.7.0 - depends_on: - kafka-1: - condition: service_healthy - kafka-2: - condition: service_healthy - kafka-3: - condition: service_healthy - ports: - - "127.0.0.1:8085:8081" - environment: - SCHEMA_REGISTRY_HOST_NAME: schema-registry - SCHEMA_REGISTRY_LISTENERS: http://0.0.0.0:8081 - SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: PLAINTEXT://kafka-1:9092,PLAINTEXT://kafka-2:9092,PLAINTEXT://kafka-3:9092 - healthcheck: - test: ["CMD-SHELL", "curl -fsS http://localhost:8081/subjects >/dev/null"] - interval: 10s - timeout: 10s - retries: 10 - - kafka-ui: - image: provectuslabs/kafka-ui:v0.7.2 - depends_on: - schema-registry: - condition: service_healthy - ports: - - "127.0.0.1:8080:8080" - environment: - DYNAMIC_CONFIG_ENABLED: "true" - KAFKA_CLUSTERS_0_NAME: agentflow-prod - KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka-1:9092,kafka-2:9092,kafka-3:9092 - KAFKA_CLUSTERS_0_SCHEMAREGISTRY: http://schema-registry:8081 - - redis: - image: redis:7.4-alpine - ports: - - "127.0.0.1:6379:6379" - command: ["redis-server", "--appendonly", "yes"] - volumes: - - redis-data:/data - healthcheck: - test: ["CMD", "redis-cli", "ping"] - interval: 10s - timeout: 5s - retries: 10 - - jaeger: - image: jaegertracing/all-in-one:1.55 - ports: - - "127.0.0.1:16686:16686" - - "127.0.0.1:4317:4317" - environment: - COLLECTOR_OTLP_ENABLED: "true" - - clickhouse: - # Part of the default bring-up since ADR 0006 fixed the serving engine on - # ClickHouse (no longer profile-gated). - image: clickhouse/clickhouse-server:24.8 - ports: - - "127.0.0.1:8123:8123" - - "127.0.0.1:9000:9000" - environment: - CLICKHOUSE_DB: agentflow - CLICKHOUSE_USER: ${CLICKHOUSE_USER:-default} - CLICKHOUSE_PASSWORD: ${CLICKHOUSE_PASSWORD:-} - CLICKHOUSE_DEFAULT_ACCESS_MANAGEMENT: "1" - volumes: - - clickhouse-data:/var/lib/clickhouse - healthcheck: - test: ["CMD-SHELL", "wget -qO- http://localhost:8123/ping | grep -q 'Ok.'"] - interval: 10s - timeout: 10s - retries: 10 - - serving-init: - # Provisioning ran on API boot until audit P0-2: the serving identity then - # needed CREATE/ALTER/INSERT, several booting replicas could race on the - # same seed, and an empty production store got demo rows just for being - # empty. It is a one-shot job now — the Compose counterpart of the Helm - # pre-install migration hook. - # - # This stack is a production-SHAPED local demo, so it also passes --seed and - # the documented demo entities exist. A real deployment runs --schema alone. - build: - context: . - dockerfile: Dockerfile.api - depends_on: - clickhouse: - condition: service_healthy - environment: - AGENTFLOW_SERVING_CONFIG: /app/config/serving.yaml - DUCKDB_PATH: /app/data/agentflow_api.duckdb - SERVING_BACKEND: ${SERVING_BACKEND:-clickhouse} - CLICKHOUSE_HOST: clickhouse - CLICKHOUSE_PORT: 8123 - CLICKHOUSE_DATABASE: agentflow - CLICKHOUSE_USER: ${CLICKHOUSE_USER:-default} - CLICKHOUSE_PASSWORD: ${CLICKHOUSE_PASSWORD:-} - command: ["python", "-m", "src.serving.provision", "--schema", "--seed"] - restart: "no" - volumes: - - agentflow-api-data:/app/data - - ./config:/app/config:ro - - agentflow-api: - build: - context: . - dockerfile: Dockerfile.api - depends_on: - jaeger: - condition: service_started - redis: - condition: service_healthy - clickhouse: - condition: service_healthy - serving-init: - condition: service_completed_successfully - ports: - - "8000:8000" - environment: - AGENTFLOW_SERVING_CONFIG: /app/config/serving.yaml - DUCKDB_PATH: /app/data/agentflow_api.duckdb - SERVING_BACKEND: ${SERVING_BACKEND:-clickhouse} - CLICKHOUSE_HOST: clickhouse - CLICKHOUSE_PORT: 8123 - CLICKHOUSE_DATABASE: agentflow - CLICKHOUSE_USER: ${CLICKHOUSE_USER:-default} - CLICKHOUSE_PASSWORD: ${CLICKHOUSE_PASSWORD:-} - OTEL_EXPORTER_OTLP_ENDPOINT: http://jaeger:4317 - OTEL_SERVICE_NAME: agentflow-api - REDIS_URL: redis://redis:6379/0 - command: ["uvicorn", "src.serving.api.main:app", "--host", "0.0.0.0", "--port", "8000"] - volumes: - - agentflow-api-data:/app/data - - ./config:/app/config:ro - healthcheck: - # /v1/health always answers 200 — its status is in the payload — so this - # check reported "healthy" for a container whose serving store was dead - # (audit P0-3). /health/ready answers 503 instead, and urlopen raises. - test: - - CMD - - python - - -c - - import urllib.request; urllib.request.urlopen("http://localhost:8000/health/ready") - interval: 10s - timeout: 10s - retries: 10 - - prometheus: - image: prom/prometheus:v2.54.0 - depends_on: - agentflow-api: - condition: service_healthy - ports: - - "127.0.0.1:9090:9090" - entrypoint: - - /bin/sh - - -ec - - | - cat <<'EOF' >/etc/prometheus/prometheus.yml - global: - scrape_interval: 15s - evaluation_interval: 15s - - scrape_configs: - - job_name: "agentflow-api" - metrics_path: /metrics - static_configs: - - targets: ["agentflow-api:8000"] - EOF - exec /bin/prometheus --config.file=/etc/prometheus/prometheus.yml --storage.tsdb.path=/prometheus --web.enable-lifecycle - volumes: - - prometheus-data:/prometheus - - grafana: - image: grafana/grafana:11.2.0 - depends_on: - prometheus: - condition: service_started - ports: - - "127.0.0.1:3000:3000" - environment: - GF_SECURITY_ADMIN_USER: ${GF_SECURITY_ADMIN_USER:-} - GF_SECURITY_ADMIN_PASSWORD: ${GF_SECURITY_ADMIN_PASSWORD:-} - GF_DASHBOARDS_DEFAULT_HOME_DASHBOARD_PATH: /var/lib/grafana/dashboards/pipeline_health.json - volumes: - - ./monitoring/grafana/dashboards:/var/lib/grafana/dashboards:ro - - ./monitoring/grafana/provisioning/datasources:/etc/grafana/provisioning/datasources:ro - - ./monitoring/grafana/provisioning/dashboards:/etc/grafana/provisioning/dashboards:ro - - grafana-data:/var/lib/grafana - -volumes: - kafka-1-data: - kafka-2-data: - kafka-3-data: - redis-data: - clickhouse-data: - agentflow-api-data: - prometheus-data: - grafana-data: +# Bounded container logs: the 2026-07-18 soak filled the stand's disk with +# unrotated json-file logs (TM alone grew to 3.8G) and took the pipeline down. +x-logging: &default-logging + driver: json-file + options: + max-size: "50m" + max-file: "3" + +services: + kafka-1: + logging: *default-logging + image: confluentinc/cp-kafka:7.7.0 + hostname: kafka-1 + ports: + - "127.0.0.1:19092:19092" + environment: + CLUSTER_ID: "AgentFlowKafkaProdCluster01" + KAFKA_NODE_ID: 1 + KAFKA_PROCESS_ROLES: broker,controller + KAFKA_LISTENERS: INTERNAL://0.0.0.0:9092,EXTERNAL://0.0.0.0:19092,CONTROLLER://0.0.0.0:29093 + KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka-1:9092,EXTERNAL://localhost:19092 + KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT,CONTROLLER:PLAINTEXT + KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL + KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER + KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka-1:29093,2@kafka-2:29093,3@kafka-3:29093 + KAFKA_LOG_DIRS: /var/lib/kafka/data + KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true" + KAFKA_DEFAULT_REPLICATION_FACTOR: 3 + KAFKA_MIN_INSYNC_REPLICAS: 2 + KAFKA_NUM_PARTITIONS: 6 + KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3 + KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3 + KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2 + KAFKA_CONFLUENT_LICENSE_TOPIC_REPLICATION_FACTOR: 3 + volumes: + - kafka-1-data:/var/lib/kafka/data + healthcheck: + test: ["CMD-SHELL", "kafka-broker-api-versions --bootstrap-server localhost:9092"] + interval: 10s + timeout: 10s + retries: 10 + + kafka-2: + logging: *default-logging + image: confluentinc/cp-kafka:7.7.0 + hostname: kafka-2 + ports: + - "127.0.0.1:29092:19092" + environment: + CLUSTER_ID: "AgentFlowKafkaProdCluster01" + KAFKA_NODE_ID: 2 + KAFKA_PROCESS_ROLES: broker,controller + KAFKA_LISTENERS: INTERNAL://0.0.0.0:9092,EXTERNAL://0.0.0.0:19092,CONTROLLER://0.0.0.0:29093 + KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka-2:9092,EXTERNAL://localhost:29092 + KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT,CONTROLLER:PLAINTEXT + KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL + KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER + KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka-1:29093,2@kafka-2:29093,3@kafka-3:29093 + KAFKA_LOG_DIRS: /var/lib/kafka/data + KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true" + KAFKA_DEFAULT_REPLICATION_FACTOR: 3 + KAFKA_MIN_INSYNC_REPLICAS: 2 + KAFKA_NUM_PARTITIONS: 6 + KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3 + KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3 + KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2 + KAFKA_CONFLUENT_LICENSE_TOPIC_REPLICATION_FACTOR: 3 + volumes: + - kafka-2-data:/var/lib/kafka/data + healthcheck: + test: ["CMD-SHELL", "kafka-broker-api-versions --bootstrap-server localhost:9092"] + interval: 10s + timeout: 10s + retries: 10 + + kafka-3: + logging: *default-logging + image: confluentinc/cp-kafka:7.7.0 + hostname: kafka-3 + ports: + - "127.0.0.1:39092:19092" + environment: + CLUSTER_ID: "AgentFlowKafkaProdCluster01" + KAFKA_NODE_ID: 3 + KAFKA_PROCESS_ROLES: broker,controller + KAFKA_LISTENERS: INTERNAL://0.0.0.0:9092,EXTERNAL://0.0.0.0:19092,CONTROLLER://0.0.0.0:29093 + KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka-3:9092,EXTERNAL://localhost:39092 + KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT,CONTROLLER:PLAINTEXT + KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL + KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER + KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka-1:29093,2@kafka-2:29093,3@kafka-3:29093 + KAFKA_LOG_DIRS: /var/lib/kafka/data + KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true" + KAFKA_DEFAULT_REPLICATION_FACTOR: 3 + KAFKA_MIN_INSYNC_REPLICAS: 2 + KAFKA_NUM_PARTITIONS: 6 + KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3 + KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3 + KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2 + KAFKA_CONFLUENT_LICENSE_TOPIC_REPLICATION_FACTOR: 3 + volumes: + - kafka-3-data:/var/lib/kafka/data + healthcheck: + test: ["CMD-SHELL", "kafka-broker-api-versions --bootstrap-server localhost:9092"] + interval: 10s + timeout: 10s + retries: 10 + + schema-registry: + logging: *default-logging + image: confluentinc/cp-schema-registry:7.7.0 + depends_on: + kafka-1: + condition: service_healthy + kafka-2: + condition: service_healthy + kafka-3: + condition: service_healthy + ports: + - "127.0.0.1:8085:8081" + environment: + SCHEMA_REGISTRY_HOST_NAME: schema-registry + SCHEMA_REGISTRY_LISTENERS: http://0.0.0.0:8081 + SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: PLAINTEXT://kafka-1:9092,PLAINTEXT://kafka-2:9092,PLAINTEXT://kafka-3:9092 + healthcheck: + test: ["CMD-SHELL", "curl -fsS http://localhost:8081/subjects >/dev/null"] + interval: 10s + timeout: 10s + retries: 10 + + kafka-ui: + logging: *default-logging + image: provectuslabs/kafka-ui:v0.7.2 + depends_on: + schema-registry: + condition: service_healthy + ports: + - "127.0.0.1:8080:8080" + environment: + DYNAMIC_CONFIG_ENABLED: "true" + KAFKA_CLUSTERS_0_NAME: agentflow-prod + KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka-1:9092,kafka-2:9092,kafka-3:9092 + KAFKA_CLUSTERS_0_SCHEMAREGISTRY: http://schema-registry:8081 + + redis: + logging: *default-logging + image: redis:7.4-alpine + ports: + - "127.0.0.1:6379:6379" + command: ["redis-server", "--appendonly", "yes"] + volumes: + - redis-data:/data + healthcheck: + test: ["CMD", "redis-cli", "ping"] + interval: 10s + timeout: 5s + retries: 10 + + jaeger: + logging: *default-logging + image: jaegertracing/all-in-one:1.55 + ports: + - "127.0.0.1:16686:16686" + - "127.0.0.1:4317:4317" + environment: + COLLECTOR_OTLP_ENABLED: "true" + + clickhouse: + logging: *default-logging + # Part of the default bring-up since ADR 0006 fixed the serving engine on + # ClickHouse (no longer profile-gated). + image: clickhouse/clickhouse-server:24.8 + ports: + - "127.0.0.1:8123:8123" + - "127.0.0.1:9000:9000" + environment: + CLICKHOUSE_DB: agentflow + CLICKHOUSE_USER: ${CLICKHOUSE_USER:-default} + CLICKHOUSE_PASSWORD: ${CLICKHOUSE_PASSWORD:-} + CLICKHOUSE_DEFAULT_ACCESS_MANAGEMENT: "1" + volumes: + - clickhouse-data:/var/lib/clickhouse + healthcheck: + test: ["CMD-SHELL", "wget -qO- http://localhost:8123/ping | grep -q 'Ok.'"] + interval: 10s + timeout: 10s + retries: 10 + + serving-init: + logging: *default-logging + # Provisioning ran on API boot until audit P0-2: the serving identity then + # needed CREATE/ALTER/INSERT, several booting replicas could race on the + # same seed, and an empty production store got demo rows just for being + # empty. It is a one-shot job now — the Compose counterpart of the Helm + # pre-install migration hook. + # + # This stack is a production-SHAPED local demo, so it also passes --seed and + # the documented demo entities exist. A real deployment runs --schema alone. + build: + context: . + dockerfile: Dockerfile.api + depends_on: + clickhouse: + condition: service_healthy + environment: + AGENTFLOW_SERVING_CONFIG: /app/config/serving.yaml + DUCKDB_PATH: /app/data/agentflow_api.duckdb + SERVING_BACKEND: ${SERVING_BACKEND:-clickhouse} + CLICKHOUSE_HOST: clickhouse + CLICKHOUSE_PORT: 8123 + CLICKHOUSE_DATABASE: agentflow + CLICKHOUSE_USER: ${CLICKHOUSE_USER:-default} + CLICKHOUSE_PASSWORD: ${CLICKHOUSE_PASSWORD:-} + command: ["python", "-m", "src.serving.provision", "--schema", "--seed"] + restart: "no" + volumes: + - agentflow-api-data:/app/data + - ./config:/app/config:ro + + agentflow-api: + logging: *default-logging + build: + context: . + dockerfile: Dockerfile.api + depends_on: + jaeger: + condition: service_started + redis: + condition: service_healthy + clickhouse: + condition: service_healthy + serving-init: + condition: service_completed_successfully + ports: + - "8000:8000" + environment: + AGENTFLOW_SERVING_CONFIG: /app/config/serving.yaml + DUCKDB_PATH: /app/data/agentflow_api.duckdb + SERVING_BACKEND: ${SERVING_BACKEND:-clickhouse} + CLICKHOUSE_HOST: clickhouse + CLICKHOUSE_PORT: 8123 + CLICKHOUSE_DATABASE: agentflow + CLICKHOUSE_USER: ${CLICKHOUSE_USER:-default} + CLICKHOUSE_PASSWORD: ${CLICKHOUSE_PASSWORD:-} + OTEL_EXPORTER_OTLP_ENDPOINT: http://jaeger:4317 + OTEL_SERVICE_NAME: agentflow-api + REDIS_URL: redis://redis:6379/0 + command: ["uvicorn", "src.serving.api.main:app", "--host", "0.0.0.0", "--port", "8000"] + volumes: + - agentflow-api-data:/app/data + - ./config:/app/config:ro + healthcheck: + # /v1/health always answers 200 — its status is in the payload — so this + # check reported "healthy" for a container whose serving store was dead + # (audit P0-3). /health/ready answers 503 instead, and urlopen raises. + test: + - CMD + - python + - -c + - import urllib.request; urllib.request.urlopen("http://localhost:8000/health/ready") + interval: 10s + timeout: 10s + retries: 10 + + prometheus: + logging: *default-logging + image: prom/prometheus:v2.54.0 + depends_on: + agentflow-api: + condition: service_healthy + ports: + - "127.0.0.1:9090:9090" + entrypoint: + - /bin/sh + - -ec + - | + cat <<'EOF' >/etc/prometheus/prometheus.yml + global: + scrape_interval: 15s + evaluation_interval: 15s + + scrape_configs: + - job_name: "agentflow-api" + metrics_path: /metrics + static_configs: + - targets: ["agentflow-api:8000"] + EOF + exec /bin/prometheus --config.file=/etc/prometheus/prometheus.yml --storage.tsdb.path=/prometheus --web.enable-lifecycle + volumes: + - prometheus-data:/prometheus + + grafana: + logging: *default-logging + image: grafana/grafana:11.2.0 + depends_on: + prometheus: + condition: service_started + ports: + - "127.0.0.1:3000:3000" + environment: + GF_SECURITY_ADMIN_USER: ${GF_SECURITY_ADMIN_USER:-} + GF_SECURITY_ADMIN_PASSWORD: ${GF_SECURITY_ADMIN_PASSWORD:-} + GF_DASHBOARDS_DEFAULT_HOME_DASHBOARD_PATH: /var/lib/grafana/dashboards/pipeline_health.json + volumes: + - ./monitoring/grafana/dashboards:/var/lib/grafana/dashboards:ro + - ./monitoring/grafana/provisioning/datasources:/etc/grafana/provisioning/datasources:ro + - ./monitoring/grafana/provisioning/dashboards:/etc/grafana/provisioning/dashboards:ro + - grafana-data:/var/lib/grafana + +volumes: + kafka-1-data: + kafka-2-data: + kafka-3-data: + redis-data: + clickhouse-data: + agentflow-api-data: + prometheus-data: + grafana-data: diff --git a/docker-compose.yml b/docker-compose.yml index 2a7c6699..a731a87a 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,223 +1,241 @@ -services: - # ── Kafka (KRaft mode, no Zookeeper) ────────────────────────── - kafka: - image: confluentinc/cp-kafka:7.7.0 - ports: - - "9092:9092" - - "9101:9101" - environment: - KAFKA_NODE_ID: 1 - KAFKA_PROCESS_ROLES: broker,controller - KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:29093 - KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:29093 - KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092 - KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER - KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT - KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 - KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 - KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 - KAFKA_LOG_RETENTION_HOURS: 24 - KAFKA_AUTO_CREATE_TOPICS_ENABLE: "false" - # Single-node KRaft heartbeats broker->controller over TCP inside one JVM; - # when the host stalls longer than the default 9s session the broker fences - # ITSELF and every partition goes leaderless. 45s absorbs the stalls a - # loaded dev host actually produces. - KAFKA_BROKER_SESSION_TIMEOUT_MS: "45000" - CLUSTER_ID: "AgentFlowKafkaCluster01" - KAFKA_JMX_PORT: 9101 - volumes: - - kafka-data:/var/lib/kafka/data - healthcheck: - test: kafka-broker-api-versions --bootstrap-server localhost:9092 - interval: 10s - timeout: 10s - retries: 5 - - kafka-init: - image: confluentinc/cp-kafka:7.7.0 - depends_on: - kafka: - condition: service_healthy - entrypoint: ["/bin/bash", "-c"] - command: - - | - set -e - echo "Creating Kafka topics..." - kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic orders.raw --partitions 6 --replication-factor 1 - kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic payments.raw --partitions 6 --replication-factor 1 - kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic clicks.raw --partitions 6 --replication-factor 1 - kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic products.cdc --partitions 3 --replication-factor 1 - kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic events.validated --partitions 6 --replication-factor 1 - kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic events.deadletter --partitions 3 --replication-factor 1 - # The stream processor's KafkaSource subscribes to the CDC topics too - # (src/processing/flink_jobs/stream_processor.py). Kafka refuses to hand - # out metadata for a topic that does not exist, so the job dies in a - # restart loop -- "Failed to get metadata for topics [...]" -- before it - # ever reads orders.raw. Debezium creates these in a full deployment; - # locally they stay empty, and empty is fine. Absent is not. - kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic cdc.postgres.public.orders_v2 --partitions 1 --replication-factor 1 - kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic cdc.postgres.public.users_enriched --partitions 1 --replication-factor 1 - kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic cdc.mysql.agentflow_demo.products_current --partitions 1 --replication-factor 1 - kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic cdc.mysql.agentflow_demo.sessions_aggregated --partitions 1 --replication-factor 1 - echo "Topics created." - - # ── Flink ───────────────────────────────────────────────────── - flink-jobmanager: - image: flink:2.3.0-java17 - ports: - - "8081:8081" - command: jobmanager - environment: - FLINK_PROPERTIES: | - jobmanager.rpc.address: flink-jobmanager - jobmanager.bind-host: 0.0.0.0 - rest.bind-host: 0.0.0.0 - state.backend.type: ${FLINK_STATE_BACKEND:-rocksdb} - execution.checkpointing.dir: s3://agentflow-lake/checkpoints - s3.endpoint: http://minio:9000 - s3.access-key: minio - s3.secret-key: minio123 - s3.path.style.access: true - execution.checkpointing.interval: 30s - execution.checkpointing.mode: EXACTLY_ONCE - volumes: - - flink-checkpoints:/opt/flink/checkpoints - healthcheck: - test: ["CMD", "curl", "-f", "http://localhost:8081/overview"] - interval: 10s - timeout: 5s - retries: 5 - - flink-taskmanager: - image: flink:2.3.0-java17 - depends_on: - flink-jobmanager: - condition: service_healthy - command: taskmanager - environment: - FLINK_PROPERTIES: | - jobmanager.rpc.address: flink-jobmanager - taskmanager.bind-host: 0.0.0.0 - taskmanager.numberOfTaskSlots: 4 - state.backend.type: ${FLINK_STATE_BACKEND:-rocksdb} - s3.endpoint: http://minio:9000 - s3.access-key: minio - s3.secret-key: minio123 - s3.path.style.access: true - deploy: - replicas: 2 - - # ── MinIO (S3-compatible object store) ──────────────────────── - minio: - image: minio/minio:RELEASE.2025-09-07T16-13-09Z - ports: - - "9000:9000" - - "9001:9001" - command: server /data --console-address ":9001" - environment: - MINIO_ROOT_USER: minio - MINIO_ROOT_PASSWORD: minio123 - volumes: - - minio-data:/data - healthcheck: - test: ["CMD", "curl", "-f", "http://localhost:9000/minio/health/live"] - interval: 10s - timeout: 5s - retries: 3 - - minio-init: - image: minio/mc:RELEASE.2025-08-13T08-35-41Z - depends_on: - minio: - condition: service_healthy - entrypoint: ["/bin/bash", "-c"] - command: - - | - set -e - # `local` is one of mc's BUILT-IN aliases (http://localhost:9000). When - # `alias set local` failed, every `mc mb local/...` below silently went - # to the mc container's own loopback, failed too, and the script still - # printed success -- leaving Flink to block forever on a checkpoint - # bucket that did not exist. Use a name mc does not ship, and let a - # failure fail the container. - for attempt in $$(seq 1 30); do - mc alias set lake http://minio:9000 minio minio123 && break - echo "waiting for minio (attempt $$attempt)" - sleep 2 - done - # `-p` IS `--ignore-existing` in mc (one flag, two spellings); passing - # both is a hard error. - mc mb -p lake/agentflow-lake - mc mb -p lake/agentflow-lake/warehouse - mc mb -p lake/agentflow-lake/checkpoints - mc ls lake/agentflow-lake - echo "MinIO buckets created." - - redis: - image: redis:7-alpine - ports: - - "6379:6379" - healthcheck: - test: ["CMD", "redis-cli", "ping"] - interval: 5s - timeout: 3s - retries: 3 - - # ── ClickHouse (serving store, ADR 0006) ────────────────────── - clickhouse: - image: clickhouse/clickhouse-server:24.8 - ports: - - "127.0.0.1:8123:8123" - # ClickHouse's native protocol, published on 9010 rather than its default - # 9000: MinIO already owns host 9000 in this file, and config/iceberg.yaml - # points the S3 endpoint at localhost:9000. Publishing both collided, so - # the Flink/MinIO lake path and the ClickHouse serving backend could not - # be brought up together — which is exactly the topology the serving - # bridge needs (Flink -> events.validated -> bridge -> ClickHouse). - # Nothing in the project speaks the native protocol; every ClickHouse - # client here uses HTTP on 8123 (config/serving.yaml, CLICKHOUSE_PORT). - - "127.0.0.1:9010:9000" - environment: - CLICKHOUSE_DB: agentflow - CLICKHOUSE_USER: ${CLICKHOUSE_USER:-default} - CLICKHOUSE_PASSWORD: ${CLICKHOUSE_PASSWORD:-} - CLICKHOUSE_DEFAULT_ACCESS_MANAGEMENT: "1" - volumes: - - clickhouse-data:/var/lib/clickhouse - healthcheck: - test: ["CMD-SHELL", "wget -qO- http://localhost:8123/ping | grep -q 'Ok.'"] - interval: 10s - timeout: 10s - retries: 10 - - # ── Prometheus ──────────────────────────────────────────────── - prometheus: - image: prom/prometheus:v2.54.0 - ports: - - "9090:9090" - volumes: - - ./monitoring/prometheus/prometheus.yml:/etc/prometheus/prometheus.yml - - ./monitoring/alerting/rules.yml:/etc/prometheus/rules.yml - - prometheus-data:/prometheus - - # ── Grafana ─────────────────────────────────────────────────── - grafana: - image: grafana/grafana:11.2.0 - ports: - - "3000:3000" - environment: - GF_SECURITY_ADMIN_PASSWORD: admin - GF_DASHBOARDS_DEFAULT_HOME_DASHBOARD_PATH: /var/lib/grafana/dashboards/pipeline_health.json - volumes: - - ./monitoring/grafana/dashboards:/var/lib/grafana/dashboards - - grafana-data:/var/lib/grafana - depends_on: - - prometheus - -volumes: - kafka-data: - flink-checkpoints: - minio-data: - clickhouse-data: - prometheus-data: - grafana-data: +# Bounded container logs: the 2026-07-18 soak filled the stand's disk with +# unrotated json-file logs (TM alone grew to 3.8G) and took the pipeline down. +x-logging: &default-logging + driver: json-file + options: + max-size: "50m" + max-file: "3" + +services: + # ── Kafka (KRaft mode, no Zookeeper) ────────────────────────── + kafka: + logging: *default-logging + image: confluentinc/cp-kafka:7.7.0 + ports: + - "9092:9092" + - "9101:9101" + environment: + KAFKA_NODE_ID: 1 + KAFKA_PROCESS_ROLES: broker,controller + KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:29093 + KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:29093 + KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092 + KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER + KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT + KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 + KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 + KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 + KAFKA_LOG_RETENTION_HOURS: 24 + KAFKA_AUTO_CREATE_TOPICS_ENABLE: "false" + # Single-node KRaft heartbeats broker->controller over TCP inside one JVM; + # when the host stalls longer than the default 9s session the broker fences + # ITSELF and every partition goes leaderless. 45s absorbs the stalls a + # loaded dev host actually produces. + KAFKA_BROKER_SESSION_TIMEOUT_MS: "45000" + CLUSTER_ID: "AgentFlowKafkaCluster01" + KAFKA_JMX_PORT: 9101 + volumes: + - kafka-data:/var/lib/kafka/data + healthcheck: + test: kafka-broker-api-versions --bootstrap-server localhost:9092 + interval: 10s + timeout: 10s + retries: 5 + + kafka-init: + logging: *default-logging + image: confluentinc/cp-kafka:7.7.0 + depends_on: + kafka: + condition: service_healthy + entrypoint: ["/bin/bash", "-c"] + command: + - | + set -e + echo "Creating Kafka topics..." + kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic orders.raw --partitions 6 --replication-factor 1 + kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic payments.raw --partitions 6 --replication-factor 1 + kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic clicks.raw --partitions 6 --replication-factor 1 + kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic products.cdc --partitions 3 --replication-factor 1 + kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic events.validated --partitions 6 --replication-factor 1 + kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic events.deadletter --partitions 3 --replication-factor 1 + # The stream processor's KafkaSource subscribes to the CDC topics too + # (src/processing/flink_jobs/stream_processor.py). Kafka refuses to hand + # out metadata for a topic that does not exist, so the job dies in a + # restart loop -- "Failed to get metadata for topics [...]" -- before it + # ever reads orders.raw. Debezium creates these in a full deployment; + # locally they stay empty, and empty is fine. Absent is not. + kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic cdc.postgres.public.orders_v2 --partitions 1 --replication-factor 1 + kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic cdc.postgres.public.users_enriched --partitions 1 --replication-factor 1 + kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic cdc.mysql.agentflow_demo.products_current --partitions 1 --replication-factor 1 + kafka-topics --bootstrap-server kafka:9092 --create --if-not-exists --topic cdc.mysql.agentflow_demo.sessions_aggregated --partitions 1 --replication-factor 1 + echo "Topics created." + + # ── Flink ───────────────────────────────────────────────────── + flink-jobmanager: + logging: *default-logging + image: flink:2.3.0-java17 + ports: + - "8081:8081" + command: jobmanager + environment: + FLINK_PROPERTIES: | + jobmanager.rpc.address: flink-jobmanager + jobmanager.bind-host: 0.0.0.0 + rest.bind-host: 0.0.0.0 + state.backend.type: ${FLINK_STATE_BACKEND:-rocksdb} + execution.checkpointing.dir: s3://agentflow-lake/checkpoints + s3.endpoint: http://minio:9000 + s3.access-key: minio + s3.secret-key: minio123 + s3.path.style.access: true + execution.checkpointing.interval: 30s + execution.checkpointing.mode: EXACTLY_ONCE + volumes: + - flink-checkpoints:/opt/flink/checkpoints + healthcheck: + test: ["CMD", "curl", "-f", "http://localhost:8081/overview"] + interval: 10s + timeout: 5s + retries: 5 + + flink-taskmanager: + logging: *default-logging + image: flink:2.3.0-java17 + depends_on: + flink-jobmanager: + condition: service_healthy + command: taskmanager + environment: + FLINK_PROPERTIES: | + jobmanager.rpc.address: flink-jobmanager + taskmanager.bind-host: 0.0.0.0 + taskmanager.numberOfTaskSlots: 4 + state.backend.type: ${FLINK_STATE_BACKEND:-rocksdb} + s3.endpoint: http://minio:9000 + s3.access-key: minio + s3.secret-key: minio123 + s3.path.style.access: true + deploy: + replicas: 2 + + # ── MinIO (S3-compatible object store) ──────────────────────── + minio: + logging: *default-logging + image: minio/minio:RELEASE.2025-09-07T16-13-09Z + ports: + - "9000:9000" + - "9001:9001" + command: server /data --console-address ":9001" + environment: + MINIO_ROOT_USER: minio + MINIO_ROOT_PASSWORD: minio123 + volumes: + - minio-data:/data + healthcheck: + test: ["CMD", "curl", "-f", "http://localhost:9000/minio/health/live"] + interval: 10s + timeout: 5s + retries: 3 + + minio-init: + logging: *default-logging + image: minio/mc:RELEASE.2025-08-13T08-35-41Z + depends_on: + minio: + condition: service_healthy + entrypoint: ["/bin/bash", "-c"] + command: + - | + set -e + # `local` is one of mc's BUILT-IN aliases (http://localhost:9000). When + # `alias set local` failed, every `mc mb local/...` below silently went + # to the mc container's own loopback, failed too, and the script still + # printed success -- leaving Flink to block forever on a checkpoint + # bucket that did not exist. Use a name mc does not ship, and let a + # failure fail the container. + for attempt in $$(seq 1 30); do + mc alias set lake http://minio:9000 minio minio123 && break + echo "waiting for minio (attempt $$attempt)" + sleep 2 + done + # `-p` IS `--ignore-existing` in mc (one flag, two spellings); passing + # both is a hard error. + mc mb -p lake/agentflow-lake + mc mb -p lake/agentflow-lake/warehouse + mc mb -p lake/agentflow-lake/checkpoints + mc ls lake/agentflow-lake + echo "MinIO buckets created." + + redis: + logging: *default-logging + image: redis:7-alpine + ports: + - "6379:6379" + healthcheck: + test: ["CMD", "redis-cli", "ping"] + interval: 5s + timeout: 3s + retries: 3 + + # ── ClickHouse (serving store, ADR 0006) ────────────────────── + clickhouse: + logging: *default-logging + image: clickhouse/clickhouse-server:24.8 + ports: + - "127.0.0.1:8123:8123" + # ClickHouse's native protocol, published on 9010 rather than its default + # 9000: MinIO already owns host 9000 in this file, and config/iceberg.yaml + # points the S3 endpoint at localhost:9000. Publishing both collided, so + # the Flink/MinIO lake path and the ClickHouse serving backend could not + # be brought up together — which is exactly the topology the serving + # bridge needs (Flink -> events.validated -> bridge -> ClickHouse). + # Nothing in the project speaks the native protocol; every ClickHouse + # client here uses HTTP on 8123 (config/serving.yaml, CLICKHOUSE_PORT). + - "127.0.0.1:9010:9000" + environment: + CLICKHOUSE_DB: agentflow + CLICKHOUSE_USER: ${CLICKHOUSE_USER:-default} + CLICKHOUSE_PASSWORD: ${CLICKHOUSE_PASSWORD:-} + CLICKHOUSE_DEFAULT_ACCESS_MANAGEMENT: "1" + volumes: + - clickhouse-data:/var/lib/clickhouse + healthcheck: + test: ["CMD-SHELL", "wget -qO- http://localhost:8123/ping | grep -q 'Ok.'"] + interval: 10s + timeout: 10s + retries: 10 + + # ── Prometheus ──────────────────────────────────────────────── + prometheus: + logging: *default-logging + image: prom/prometheus:v2.54.0 + ports: + - "9090:9090" + volumes: + - ./monitoring/prometheus/prometheus.yml:/etc/prometheus/prometheus.yml + - ./monitoring/alerting/rules.yml:/etc/prometheus/rules.yml + - prometheus-data:/prometheus + + # ── Grafana ─────────────────────────────────────────────────── + grafana: + logging: *default-logging + image: grafana/grafana:11.2.0 + ports: + - "3000:3000" + environment: + GF_SECURITY_ADMIN_PASSWORD: admin + GF_DASHBOARDS_DEFAULT_HOME_DASHBOARD_PATH: /var/lib/grafana/dashboards/pipeline_health.json + volumes: + - ./monitoring/grafana/dashboards:/var/lib/grafana/dashboards + - grafana-data:/var/lib/grafana + depends_on: + - prometheus + +volumes: + kafka-data: + flink-checkpoints: + minio-data: + clickhouse-data: + prometheus-data: + grafana-data: diff --git a/src/processing/flink_jobs/session_aggregator.py b/src/processing/flink_jobs/session_aggregator.py index 2507f0fb..4ed90a77 100644 --- a/src/processing/flink_jobs/session_aggregator.py +++ b/src/processing/flink_jobs/session_aggregator.py @@ -14,6 +14,7 @@ from typing import Any from pyflink.common import Types, WatermarkStrategy +from pyflink.common.restart_strategy import RestartStrategies from pyflink.common.serialization import SimpleStringSchema from pyflink.common.time import Duration from pyflink.common.watermark_strategy import TimestampAssigner @@ -158,6 +159,16 @@ def build_pipeline() -> StreamExecutionEnvironment: env.enable_checkpointing(CHECKPOINT_INTERVAL_MS) env.set_parallelism(int(os.getenv("FLINK_PARALLELISM", "2"))) + # Bounded restart budget — same rationale as stream_processor.build_pipeline: + # a persistently failing job must park in FAILED, not restart forever. + env.set_restart_strategy( + RestartStrategies.failure_rate_restart( + int(os.getenv("FLINK_RESTART_MAX_FAILURES_PER_INTERVAL", "3")), + int(os.getenv("FLINK_RESTART_FAILURE_RATE_INTERVAL_MS", "300000")), # 5 min + int(os.getenv("FLINK_RESTART_DELAY_MS", "10000")), # 10s + ) + ) + bootstrap_servers = os.getenv("KAFKA_BOOTSTRAP_SERVERS", "localhost:9092") source = ( diff --git a/src/processing/flink_jobs/stream_processor.py b/src/processing/flink_jobs/stream_processor.py index c046cd00..ee54e8e3 100644 --- a/src/processing/flink_jobs/stream_processor.py +++ b/src/processing/flink_jobs/stream_processor.py @@ -13,6 +13,7 @@ from typing import Any from pyflink.common import Types, WatermarkStrategy +from pyflink.common.restart_strategy import RestartStrategies from pyflink.common.serialization import SimpleStringSchema from pyflink.common.time import Duration, Time from pyflink.common.watermark_strategy import TimestampAssigner @@ -260,6 +261,19 @@ def build_pipeline() -> StreamExecutionEnvironment: env.get_checkpoint_config().set_min_pause_between_checkpoints(10_000) env.set_parallelism(int(os.getenv("FLINK_PARALLELISM", "2"))) + # Bounded restart budget: with the default infinite fixed-delay strategy a + # job whose environment is persistently broken (e.g. checkpoint storage + # full) restarts forever, re-emitting from the last checkpoint on every + # attempt and flooding downstream with duplicates. Exceeding the failure + # rate parks the job in FAILED, where alerting and the operator take over. + env.set_restart_strategy( + RestartStrategies.failure_rate_restart( + int(os.getenv("FLINK_RESTART_MAX_FAILURES_PER_INTERVAL", "3")), + int(os.getenv("FLINK_RESTART_FAILURE_RATE_INTERVAL_MS", "300000")), # 5 min + int(os.getenv("FLINK_RESTART_DELAY_MS", "10000")), # 10s + ) + ) + bootstrap_servers = os.getenv("KAFKA_BOOTSTRAP_SERVERS", "localhost:9092") # Multi-topic Kafka source diff --git a/tests/unit/test_session_aggregator.py b/tests/unit/test_session_aggregator.py index 8fc604bc..de01c817 100644 --- a/tests/unit/test_session_aggregator.py +++ b/tests/unit/test_session_aggregator.py @@ -92,6 +92,7 @@ def __init__(self): self.checkpointing = None self.parallelism = None self.from_source_args = None + self.restart_strategy = None self.stream = _FakeStream() def enable_checkpointing(self, interval): @@ -100,6 +101,9 @@ def enable_checkpointing(self, interval): def set_parallelism(self, value): self.parallelism = value + def set_restart_strategy(self, value): + self.restart_strategy = value + def from_source(self, source, watermark_strategy, name): self.from_source_args = (source, watermark_strategy, name) return self.stream @@ -172,6 +176,15 @@ def of_seconds(cls, seconds): time_module.Duration = _Duration + restart_strategy = types.ModuleType("pyflink.common.restart_strategy") + + class _RestartStrategies: + @staticmethod + def failure_rate_restart(max_failure_rate, failure_rate_interval, delay_interval): + return ("failure_rate", max_failure_rate, failure_rate_interval, delay_interval) + + restart_strategy.RestartStrategies = _RestartStrategies + datastream = types.ModuleType("pyflink.datastream") datastream.__path__ = [] datastream.StreamExecutionEnvironment = _FakeExecutionEnvironment @@ -267,6 +280,7 @@ def __init__(self, name, type_info): common.serialization = serialization common.time = time_module common.watermark_strategy = watermark_strategy + common.restart_strategy = restart_strategy datastream.connectors = connectors connectors.kafka = kafka datastream.functions = functions @@ -279,6 +293,7 @@ def __init__(self, name, type_info): monkeypatch.setitem(sys.modules, "pyflink.common.serialization", serialization) monkeypatch.setitem(sys.modules, "pyflink.common.time", time_module) monkeypatch.setitem(sys.modules, "pyflink.common.watermark_strategy", watermark_strategy) + monkeypatch.setitem(sys.modules, "pyflink.common.restart_strategy", restart_strategy) monkeypatch.setitem(sys.modules, "pyflink.datastream", datastream) monkeypatch.setitem(sys.modules, "pyflink.datastream.connectors", connectors) monkeypatch.setitem(sys.modules, "pyflink.datastream.connectors.kafka", kafka) @@ -596,6 +611,21 @@ def test_build_pipeline_uses_defaults_and_wires_stream(session_aggregator, monke assert env.stream.sink["record_serializer"]["topic"] == "sessions.aggregated" +def test_build_pipeline_sets_bounded_restart_strategy(session_aggregator, monkeypatch): + env = _FakeExecutionEnvironment() + session_aggregator.StreamExecutionEnvironment.current_env = env + for name in ( + "FLINK_RESTART_MAX_FAILURES_PER_INTERVAL", + "FLINK_RESTART_FAILURE_RATE_INTERVAL_MS", + "FLINK_RESTART_DELAY_MS", + ): + monkeypatch.delenv(name, raising=False) + + session_aggregator.build_pipeline() + + assert env.restart_strategy == ("failure_rate", 3, 300_000, 10_000) + + def test_build_pipeline_respects_environment_overrides(session_aggregator, monkeypatch): env = _FakeExecutionEnvironment() session_aggregator.StreamExecutionEnvironment.current_env = env diff --git a/tests/unit/test_stream_processor.py b/tests/unit/test_stream_processor.py index ec32932d..fccb3f18 100644 --- a/tests/unit/test_stream_processor.py +++ b/tests/unit/test_stream_processor.py @@ -116,6 +116,7 @@ def __init__(self): self.parallelism = None self.checkpoint_config = _FakeCheckpointConfig() self.from_source_args = None + self.restart_strategy = None self.source_stream = _FakeSourceStream(self) self.validated_stream = None self.dead_letter_stream = None @@ -132,6 +133,9 @@ def get_checkpoint_config(self): def set_parallelism(self, value): self.parallelism = value + def set_restart_strategy(self, value): + self.restart_strategy = value + def from_source(self, source, watermark_strategy, name): self.from_source_args = (source, watermark_strategy, name) return self.source_stream @@ -300,6 +304,15 @@ def minutes(cls, minutes): time_module.Duration = _Duration time_module.Time = _Time + restart_strategy = types.ModuleType("pyflink.common.restart_strategy") + + class _RestartStrategies: + @staticmethod + def failure_rate_restart(max_failure_rate, failure_rate_interval, delay_interval): + return ("failure_rate", max_failure_rate, failure_rate_interval, delay_interval) + + restart_strategy.RestartStrategies = _RestartStrategies + datastream = types.ModuleType("pyflink.datastream") datastream.__path__ = [] datastream.StreamExecutionEnvironment = _FakeExecutionEnvironment @@ -428,6 +441,7 @@ def enable_time_to_live(self, ttl_config): common.serialization = serialization common.time = time_module common.watermark_strategy = watermark_strategy + common.restart_strategy = restart_strategy datastream.connectors = connectors connectors.kafka = kafka datastream.functions = functions @@ -441,6 +455,7 @@ def enable_time_to_live(self, ttl_config): monkeypatch.setitem(sys.modules, "pyflink.common.serialization", serialization) monkeypatch.setitem(sys.modules, "pyflink.common.time", time_module) monkeypatch.setitem(sys.modules, "pyflink.common.watermark_strategy", watermark_strategy) + monkeypatch.setitem(sys.modules, "pyflink.common.restart_strategy", restart_strategy) monkeypatch.setitem(sys.modules, "pyflink.datastream", datastream) monkeypatch.setitem(sys.modules, "pyflink.datastream.connectors", connectors) monkeypatch.setitem(sys.modules, "pyflink.datastream.connectors.kafka", kafka) @@ -852,6 +867,33 @@ def test_build_pipeline_uses_defaults_and_wires_sinks(stream_processor, monkeypa assert env.filtered_stream.sink["record_serializer"]["topic"] == "events.validated" +def test_build_pipeline_sets_bounded_restart_strategy(stream_processor, monkeypatch): + env = _FakeExecutionEnvironment() + stream_processor.StreamExecutionEnvironment.current_env = env + for name in ( + "FLINK_RESTART_MAX_FAILURES_PER_INTERVAL", + "FLINK_RESTART_FAILURE_RATE_INTERVAL_MS", + "FLINK_RESTART_DELAY_MS", + ): + monkeypatch.delenv(name, raising=False) + + stream_processor.build_pipeline() + + assert env.restart_strategy == ("failure_rate", 3, 300_000, 10_000) + + +def test_build_pipeline_restart_strategy_env_overrides(stream_processor, monkeypatch): + env = _FakeExecutionEnvironment() + stream_processor.StreamExecutionEnvironment.current_env = env + monkeypatch.setenv("FLINK_RESTART_MAX_FAILURES_PER_INTERVAL", "5") + monkeypatch.setenv("FLINK_RESTART_FAILURE_RATE_INTERVAL_MS", "600000") + monkeypatch.setenv("FLINK_RESTART_DELAY_MS", "20000") + + stream_processor.build_pipeline() + + assert env.restart_strategy == ("failure_rate", 5, 600_000, 20_000) + + def test_build_pipeline_respects_environment_overrides(stream_processor, monkeypatch): env = _FakeExecutionEnvironment() stream_processor.StreamExecutionEnvironment.current_env = env From 0b7a1c7e5db3c378076bcdd5d9a1c7de36f71fd4 Mon Sep 17 00:00:00 2001 From: JuliaEdom Date: Mon, 20 Jul 2026 16:44:04 +0300 Subject: [PATCH 2/2] fix(flink): route restart strategy through Configuration/env.configure for Flink 2.x StreamExecutionEnvironment.set_restart_strategy was removed in Flink 2.x (FLIP-381); on the pinned apache-flink 2.3.0 runtime build_pipeline() raised AttributeError at submission. Both jobs now set the failure-rate budget via Configuration.set_string + env.configure (verified against release-2.3.0 sources: Configuration re-exported in pyflink.common, env.configure at stream_execution_environment.py:379). Env overrides keep the same variables, now rendered as " ms" duration strings. Test fakes drop set_restart_strategy and the RestartStrategies module entirely so a regression back to the removed API fails in unit tests, not only on the live cluster; fixtures fake Configuration and capture env.configure instead. Co-Authored-By: Claude Fable 5 --- .../flink_jobs/session_aggregator.py | 25 +++++++---- src/processing/flink_jobs/stream_processor.py | 25 +++++++---- tests/unit/test_session_aggregator.py | 35 ++++++++------- tests/unit/test_stream_processor.py | 43 ++++++++++++------- 4 files changed, 81 insertions(+), 47 deletions(-) diff --git a/src/processing/flink_jobs/session_aggregator.py b/src/processing/flink_jobs/session_aggregator.py index 4ed90a77..f01b2a1c 100644 --- a/src/processing/flink_jobs/session_aggregator.py +++ b/src/processing/flink_jobs/session_aggregator.py @@ -13,8 +13,7 @@ from datetime import UTC, datetime from typing import Any -from pyflink.common import Types, WatermarkStrategy -from pyflink.common.restart_strategy import RestartStrategies +from pyflink.common import Configuration, Types, WatermarkStrategy from pyflink.common.serialization import SimpleStringSchema from pyflink.common.time import Duration from pyflink.common.watermark_strategy import TimestampAssigner @@ -161,13 +160,23 @@ def build_pipeline() -> StreamExecutionEnvironment: # Bounded restart budget — same rationale as stream_processor.build_pipeline: # a persistently failing job must park in FAILED, not restart forever. - env.set_restart_strategy( - RestartStrategies.failure_rate_restart( - int(os.getenv("FLINK_RESTART_MAX_FAILURES_PER_INTERVAL", "3")), - int(os.getenv("FLINK_RESTART_FAILURE_RATE_INTERVAL_MS", "300000")), # 5 min - int(os.getenv("FLINK_RESTART_DELAY_MS", "10000")), # 10s - ) + # Flink 2.x removed env.set_restart_strategy (FLIP-381); the strategy must + # go through Configuration + env.configure. + restart_config = Configuration() + restart_config.set_string("restart-strategy.type", "failure-rate") + restart_config.set_string( + "restart-strategy.failure-rate.max-failures-per-interval", + str(int(os.getenv("FLINK_RESTART_MAX_FAILURES_PER_INTERVAL", "3"))), + ) + restart_config.set_string( + "restart-strategy.failure-rate.failure-rate-interval", + f"{int(os.getenv('FLINK_RESTART_FAILURE_RATE_INTERVAL_MS', '300000'))} ms", # 5 min + ) + restart_config.set_string( + "restart-strategy.failure-rate.delay", + f"{int(os.getenv('FLINK_RESTART_DELAY_MS', '10000'))} ms", # 10s ) + env.configure(restart_config) bootstrap_servers = os.getenv("KAFKA_BOOTSTRAP_SERVERS", "localhost:9092") diff --git a/src/processing/flink_jobs/stream_processor.py b/src/processing/flink_jobs/stream_processor.py index ee54e8e3..0a14f551 100644 --- a/src/processing/flink_jobs/stream_processor.py +++ b/src/processing/flink_jobs/stream_processor.py @@ -12,8 +12,7 @@ from collections.abc import Iterator from typing import Any -from pyflink.common import Types, WatermarkStrategy -from pyflink.common.restart_strategy import RestartStrategies +from pyflink.common import Configuration, Types, WatermarkStrategy from pyflink.common.serialization import SimpleStringSchema from pyflink.common.time import Duration, Time from pyflink.common.watermark_strategy import TimestampAssigner @@ -266,13 +265,23 @@ def build_pipeline() -> StreamExecutionEnvironment: # full) restarts forever, re-emitting from the last checkpoint on every # attempt and flooding downstream with duplicates. Exceeding the failure # rate parks the job in FAILED, where alerting and the operator take over. - env.set_restart_strategy( - RestartStrategies.failure_rate_restart( - int(os.getenv("FLINK_RESTART_MAX_FAILURES_PER_INTERVAL", "3")), - int(os.getenv("FLINK_RESTART_FAILURE_RATE_INTERVAL_MS", "300000")), # 5 min - int(os.getenv("FLINK_RESTART_DELAY_MS", "10000")), # 10s - ) + # Flink 2.x removed env.set_restart_strategy (FLIP-381); the strategy must + # go through Configuration + env.configure. + restart_config = Configuration() + restart_config.set_string("restart-strategy.type", "failure-rate") + restart_config.set_string( + "restart-strategy.failure-rate.max-failures-per-interval", + str(int(os.getenv("FLINK_RESTART_MAX_FAILURES_PER_INTERVAL", "3"))), + ) + restart_config.set_string( + "restart-strategy.failure-rate.failure-rate-interval", + f"{int(os.getenv('FLINK_RESTART_FAILURE_RATE_INTERVAL_MS', '300000'))} ms", # 5 min + ) + restart_config.set_string( + "restart-strategy.failure-rate.delay", + f"{int(os.getenv('FLINK_RESTART_DELAY_MS', '10000'))} ms", # 10s ) + env.configure(restart_config) bootstrap_servers = os.getenv("KAFKA_BOOTSTRAP_SERVERS", "localhost:9092") diff --git a/tests/unit/test_session_aggregator.py b/tests/unit/test_session_aggregator.py index de01c817..17b2f7fc 100644 --- a/tests/unit/test_session_aggregator.py +++ b/tests/unit/test_session_aggregator.py @@ -92,7 +92,7 @@ def __init__(self): self.checkpointing = None self.parallelism = None self.from_source_args = None - self.restart_strategy = None + self.configured = None self.stream = _FakeStream() def enable_checkpointing(self, interval): @@ -101,8 +101,10 @@ def enable_checkpointing(self, interval): def set_parallelism(self, value): self.parallelism = value - def set_restart_strategy(self, value): - self.restart_strategy = value + # Deliberately no set_restart_strategy: Flink 2.x removed it (FLIP-381), + # so a regression back to it must fail here, not only on the live cluster. + def configure(self, configuration): + self.configured = configuration def from_source(self, source, watermark_strategy, name): self.from_source_args = (source, watermark_strategy, name) @@ -147,8 +149,16 @@ def with_timestamp_assigner(self, assigner): self.timestamp_assigner = assigner return self + class _Configuration: + def __init__(self): + self.values = {} + + def set_string(self, key, value): + self.values[key] = value + common.Types = _Types common.WatermarkStrategy = _WatermarkStrategy + common.Configuration = _Configuration serialization = types.ModuleType("pyflink.common.serialization") @@ -176,15 +186,6 @@ def of_seconds(cls, seconds): time_module.Duration = _Duration - restart_strategy = types.ModuleType("pyflink.common.restart_strategy") - - class _RestartStrategies: - @staticmethod - def failure_rate_restart(max_failure_rate, failure_rate_interval, delay_interval): - return ("failure_rate", max_failure_rate, failure_rate_interval, delay_interval) - - restart_strategy.RestartStrategies = _RestartStrategies - datastream = types.ModuleType("pyflink.datastream") datastream.__path__ = [] datastream.StreamExecutionEnvironment = _FakeExecutionEnvironment @@ -280,7 +281,6 @@ def __init__(self, name, type_info): common.serialization = serialization common.time = time_module common.watermark_strategy = watermark_strategy - common.restart_strategy = restart_strategy datastream.connectors = connectors connectors.kafka = kafka datastream.functions = functions @@ -293,7 +293,6 @@ def __init__(self, name, type_info): monkeypatch.setitem(sys.modules, "pyflink.common.serialization", serialization) monkeypatch.setitem(sys.modules, "pyflink.common.time", time_module) monkeypatch.setitem(sys.modules, "pyflink.common.watermark_strategy", watermark_strategy) - monkeypatch.setitem(sys.modules, "pyflink.common.restart_strategy", restart_strategy) monkeypatch.setitem(sys.modules, "pyflink.datastream", datastream) monkeypatch.setitem(sys.modules, "pyflink.datastream.connectors", connectors) monkeypatch.setitem(sys.modules, "pyflink.datastream.connectors.kafka", kafka) @@ -623,7 +622,13 @@ def test_build_pipeline_sets_bounded_restart_strategy(session_aggregator, monkey session_aggregator.build_pipeline() - assert env.restart_strategy == ("failure_rate", 3, 300_000, 10_000) + assert env.configured is not None + assert env.configured.values == { + "restart-strategy.type": "failure-rate", + "restart-strategy.failure-rate.max-failures-per-interval": "3", + "restart-strategy.failure-rate.failure-rate-interval": "300000 ms", + "restart-strategy.failure-rate.delay": "10000 ms", + } def test_build_pipeline_respects_environment_overrides(session_aggregator, monkeypatch): diff --git a/tests/unit/test_stream_processor.py b/tests/unit/test_stream_processor.py index fccb3f18..eb0f93d2 100644 --- a/tests/unit/test_stream_processor.py +++ b/tests/unit/test_stream_processor.py @@ -116,7 +116,7 @@ def __init__(self): self.parallelism = None self.checkpoint_config = _FakeCheckpointConfig() self.from_source_args = None - self.restart_strategy = None + self.configured = None self.source_stream = _FakeSourceStream(self) self.validated_stream = None self.dead_letter_stream = None @@ -133,8 +133,10 @@ def get_checkpoint_config(self): def set_parallelism(self, value): self.parallelism = value - def set_restart_strategy(self, value): - self.restart_strategy = value + # Deliberately no set_restart_strategy: Flink 2.x removed it (FLIP-381), + # so a regression back to it must fail here, not only on the live cluster. + def configure(self, configuration): + self.configured = configuration def from_source(self, source, watermark_strategy, name): self.from_source_args = (source, watermark_strategy, name) @@ -266,8 +268,16 @@ def with_timestamp_assigner(self, assigner): self.timestamp_assigner = assigner return self + class _Configuration: + def __init__(self): + self.values = {} + + def set_string(self, key, value): + self.values[key] = value + common.Types = _Types common.WatermarkStrategy = _WatermarkStrategy + common.Configuration = _Configuration serialization = types.ModuleType("pyflink.common.serialization") @@ -304,15 +314,6 @@ def minutes(cls, minutes): time_module.Duration = _Duration time_module.Time = _Time - restart_strategy = types.ModuleType("pyflink.common.restart_strategy") - - class _RestartStrategies: - @staticmethod - def failure_rate_restart(max_failure_rate, failure_rate_interval, delay_interval): - return ("failure_rate", max_failure_rate, failure_rate_interval, delay_interval) - - restart_strategy.RestartStrategies = _RestartStrategies - datastream = types.ModuleType("pyflink.datastream") datastream.__path__ = [] datastream.StreamExecutionEnvironment = _FakeExecutionEnvironment @@ -441,7 +442,6 @@ def enable_time_to_live(self, ttl_config): common.serialization = serialization common.time = time_module common.watermark_strategy = watermark_strategy - common.restart_strategy = restart_strategy datastream.connectors = connectors connectors.kafka = kafka datastream.functions = functions @@ -455,7 +455,6 @@ def enable_time_to_live(self, ttl_config): monkeypatch.setitem(sys.modules, "pyflink.common.serialization", serialization) monkeypatch.setitem(sys.modules, "pyflink.common.time", time_module) monkeypatch.setitem(sys.modules, "pyflink.common.watermark_strategy", watermark_strategy) - monkeypatch.setitem(sys.modules, "pyflink.common.restart_strategy", restart_strategy) monkeypatch.setitem(sys.modules, "pyflink.datastream", datastream) monkeypatch.setitem(sys.modules, "pyflink.datastream.connectors", connectors) monkeypatch.setitem(sys.modules, "pyflink.datastream.connectors.kafka", kafka) @@ -879,7 +878,13 @@ def test_build_pipeline_sets_bounded_restart_strategy(stream_processor, monkeypa stream_processor.build_pipeline() - assert env.restart_strategy == ("failure_rate", 3, 300_000, 10_000) + assert env.configured is not None + assert env.configured.values == { + "restart-strategy.type": "failure-rate", + "restart-strategy.failure-rate.max-failures-per-interval": "3", + "restart-strategy.failure-rate.failure-rate-interval": "300000 ms", + "restart-strategy.failure-rate.delay": "10000 ms", + } def test_build_pipeline_restart_strategy_env_overrides(stream_processor, monkeypatch): @@ -891,7 +896,13 @@ def test_build_pipeline_restart_strategy_env_overrides(stream_processor, monkeyp stream_processor.build_pipeline() - assert env.restart_strategy == ("failure_rate", 5, 600_000, 20_000) + assert env.configured is not None + assert env.configured.values == { + "restart-strategy.type": "failure-rate", + "restart-strategy.failure-rate.max-failures-per-interval": "5", + "restart-strategy.failure-rate.failure-rate-interval": "600000 ms", + "restart-strategy.failure-rate.delay": "20000 ms", + } def test_build_pipeline_respects_environment_overrides(stream_processor, monkeypatch):