From fcb6fa5d5d8d9e9444bfea8bc4bf06283fe5490b Mon Sep 17 00:00:00 2001 From: Tim Paine <3105306+timkpaine@users.noreply.github.com> Date: Thu, 13 Aug 2026 19:09:45 -0400 Subject: [PATCH 1/2] Setup CI for adapter with service dependencies Run the Kafka adapter integration tests in CI against a real broker. Adds a test_adapters job that stands up ci/kafka/docker-compose.yml, sets CSP_TEST_, and runs the matching tests, plus dockerup/dockerps/dockerdown targets for doing the same locally. The compose stack is trimmed to zookeeper and a single broker, with healthchecks so `docker compose up --wait` blocks until the broker accepts connections rather than relying on a fixed sleep. Only 9092 is published, bound to loopback. Broker-side topic auto-creation is disabled so test_invalid_topic can exercise the failure path, and tests create their topics explicitly through the Kafka AdminClient. Test changes target the startup race where a subscriber misses the first few records while its consumer group is being assigned. Rather than loosening the assertions, the affected tests align on the first record the subscriber saw and then require an exact contiguous run, so loss, duplication and reordering are still caught. Also drops curl from the Windows chocolatey install: the package fails whenever a new version is approved on the community feed before it is downloadable, and curl.exe ships with Windows. Signed-off-by: Tim Paine <3105306+timkpaine@users.noreply.github.com> --- .github/workflows/build.yml | 73 ++++++- Makefile | 18 +- ci/kafka/docker-compose.yml | 205 ++++--------------- conda/dev-environment-unix.yml | 2 +- conda/dev-environment-win.yml | 2 +- cpp/csp/engine/Struct.h | 2 +- csp/adapters/kafka.py | 13 +- csp/tests/adapters/conftest.py | 9 +- csp/tests/adapters/kafka_utils.py | 29 +-- csp/tests/adapters/test_kafka.py | 117 +++++------ csp/tests/adapters/test_status.py | 13 +- examples/03_using_adapters/kafka/e1_kafka.py | 6 +- pyproject.toml | 31 +-- 13 files changed, 234 insertions(+), 286 deletions(-) diff --git a/.github/workflows/build.yml b/.github/workflows/build.yml index de21a2fb9..753e3bc28 100644 --- a/.github/workflows/build.yml +++ b/.github/workflows/build.yml @@ -601,7 +601,6 @@ jobs: env: CSP_TEST_SKIP_EXAMPLES: "1" - #################################################### #..................................................# #..|########|..|########|..../####\....|########|..# @@ -713,7 +712,75 @@ jobs: ########################################################################################################### # Test Service Adapters # ########################################################################################################### - # Coming soon! + test_adapters: + needs: + - initialize + - build + + permissions: + contents: read + + timeout-minutes: 30 + + strategy: + matrix: + os: + - ubuntu-24.04 + python-version: + - 3.11 + adapter: + - kafka + + runs-on: ${{ matrix.os }} + + steps: + - name: Checkout + uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 + with: + submodules: recursive + persist-credentials: false + + - name: Set up Python ${{ matrix.python-version }} + uses: ./.github/actions/setup-python + with: + version: '${{ matrix.python-version }}' + cibuildwheel: false + + - name: Install python dependencies + run: make requirements + + - name: Install test dependencies + shell: bash + run: sudo apt-get update && sudo apt-get install -y graphviz + + # Download artifact + - name: Download wheel + uses: actions/download-artifact@3e5f45b2cfb9172054b4087a40e8e0b5a5461e7c # v8.0.1 + with: + name: csp-dist-${{ runner.os }}-${{ runner.arch }}-${{ matrix.python-version }} + + - name: Install wheel + run: | + python -m pip install -U *manylinux*.whl + python -m pip install -U --no-deps *manylinux*.whl --target . + + # Services declare healthchecks, so --wait blocks until the broker actually accepts connections + - name: Spin up adapter service + run: make dockerup ADAPTER=${{ matrix.adapter }} DOCKERARGS="--wait --wait-timeout 180" + + # Run tests + - name: Setup test flags + shell: bash + env: + ADAPTER: ${{ matrix.adapter }} + run: echo "CSP_TEST_${ADAPTER^^}=1" >> "$GITHUB_ENV" + + - name: Python Test Steps + run: make test-py TEST_ARGS="-k ${{ matrix.adapter }}" + + - name: Spin down adapter service + run: make dockerdown ADAPTER=${{ matrix.adapter }} + if: ${{ always() }} ############################################################################################ #..........................................................................................# @@ -728,7 +795,6 @@ jobs: ############################################################################################ # Upload Release Artifacts # ############################################################################################ - # only publish artifacts on tags, but otherwise this always runs # Note this whole workflow only triggers on release tags (e.g. "v0.1.0") publish_release_artifacts: @@ -740,6 +806,7 @@ jobs: - test - test_sdist - test_dependencies + - test_adapters if: startsWith(github.ref, 'refs/tags/v') runs-on: ubuntu-24.04 diff --git a/Makefile b/Makefile index b55b52d2e..4f5d49a5a 100644 --- a/Makefile +++ b/Makefile @@ -117,7 +117,9 @@ tests: test .PHONY: dockerup dockerps dockerdown initpodmanmac ADAPTER := kafka -DOCKER := podman +# Prefer docker, fall back to podman-compose; override with DOCKER_COMPOSE=... +DOCKER_COMPOSE := $(shell command -v docker >/dev/null 2>&1 && echo "docker compose" || echo "podman-compose") +DOCKERARGS := initpodmanmac: podman machine stop @@ -125,13 +127,13 @@ initpodmanmac: podman machine start dockerup: ## spin up docker compose services for adapter testing - $(DOCKER) compose -f ci/$(ADAPTER)/docker-compose.yml up -d + $(DOCKER_COMPOSE) -f ci/$(ADAPTER)/docker-compose.yml up -d $(DOCKERARGS) -dockerps: ## spin up docker compose services for adapter testing - $(DOCKER) compose -f ci/$(ADAPTER)/docker-compose.yml ps +dockerps: ## get status of current docker compose services + $(DOCKER_COMPOSE) -f ci/$(ADAPTER)/docker-compose.yml ps -dockerdown: ## spin up docker compose services for adapter testing - $(DOCKER) compose -f ci/$(ADAPTER)/docker-compose.yml down +dockerdown: ## spin down docker compose services for adapter testing + $(DOCKER_COMPOSE) -f ci/$(ADAPTER)/docker-compose.yml down ########### # VERSION # @@ -222,9 +224,11 @@ dependencies-fedora: ## install dependencies for linux - note that zip is neede dependencies-vcpkg: ## install dependencies via vcpkg cd vcpkg && ./bootstrap-vcpkg.sh && ./vcpkg install +# curl.exe ships with Windows; installing it from choco breaks whenever a new version is +# approved on the feed before the package is actually downloadable dependencies-win: ## install dependencies via windows choco install cmake --version=3.31.6 --allow-downgrade - choco install curl winflexbison ninja unzip --no-progress -y + choco install winflexbison ninja unzip --no-progress -y ############################################################################################ # Thanks to Francoise at marmelab.com for this diff --git a/ci/kafka/docker-compose.yml b/ci/kafka/docker-compose.yml index d2674945f..a15cbcce0 100644 --- a/ci/kafka/docker-compose.yml +++ b/ci/kafka/docker-compose.yml @@ -1,181 +1,48 @@ -# https://docs.confluent.io/platform/current/platform-quickstart.html -# https://raw.githubusercontent.com/confluentinc/cp-all-in-one/7.5.3-post/cp-all-in-one-kraft/docker-compose.yml +# https://github.com/conduktor/kafka-stack-docker-compose --- -version: '2' services: - zookeeper: + zoo1: image: confluentinc/cp-zookeeper:7.5.3 - hostname: zookeeper - container_name: zookeeper + hostname: zoo1 + container_name: zoo1 ports: - - "2181:2181" + - "127.0.0.1:2181:2181" environment: ZOOKEEPER_CLIENT_PORT: 2181 - ZOOKEEPER_TICK_TIME: 2000 - - broker: - image: confluentinc/cp-server:7.5.3 - hostname: broker - container_name: broker - depends_on: - - zookeeper + ZOOKEEPER_SERVER_ID: 1 + ZOOKEEPER_SERVERS: zoo1:2888:3888 + healthcheck: + # cub is shipped in the image; the 4lw "ruok" probe is not whitelisted by default in ZK 3.5+ + test: ["CMD", "cub", "zk-ready", "localhost:2181", "5"] + interval: 5s + timeout: 10s + retries: 24 + + kafka1: + image: confluentinc/cp-kafka:7.5.3 + hostname: kafka1 + container_name: kafka1 ports: - - "9092:9092" - - "9101:9101" + - "127.0.0.1:9092:9092" environment: + KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka1:19092,EXTERNAL://${DOCKER_HOST_IP:-127.0.0.1}:9092 + KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT + KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL + KAFKA_ZOOKEEPER_CONNECT: "zoo1:2181" KAFKA_BROKER_ID: 1 - KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181' - KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT - KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://broker:29092,PLAINTEXT_HOST://localhost:9092 - KAFKA_METRIC_REPORTERS: io.confluent.metrics.reporter.ConfluentMetricsReporter + KAFKA_LOG4J_LOGGERS: "kafka.controller=INFO,kafka.producer.async.DefaultEventHandler=INFO,state.change.logger=INFO" KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 - KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0 - KAFKA_CONFLUENT_LICENSE_TOPIC_REPLICATION_FACTOR: 1 - KAFKA_CONFLUENT_BALANCER_TOPIC_REPLICATION_FACTOR: 1 - KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 - KAFKA_JMX_PORT: 9101 - KAFKA_JMX_HOSTNAME: localhost - KAFKA_CONFLUENT_SCHEMA_REGISTRY_URL: http://schema-registry:8081 - CONFLUENT_METRICS_REPORTER_BOOTSTRAP_SERVERS: broker:29092 - CONFLUENT_METRICS_REPORTER_TOPIC_REPLICAS: 1 - CONFLUENT_METRICS_ENABLE: 'true' - CONFLUENT_SUPPORT_CUSTOMER_ID: 'anonymous' - KAFKA_AUTO_CREATE_TOPICS_ENABLE: 'true' - TOPIC_AUTO_CREATE: 'true' - - schema-registry: - image: confluentinc/cp-schema-registry:7.5.3 - hostname: schema-registry - container_name: schema-registry - depends_on: - - broker - ports: - - "8081:8081" - environment: - SCHEMA_REGISTRY_HOST_NAME: schema-registry - SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: 'broker:29092' - SCHEMA_REGISTRY_LISTENERS: http://0.0.0.0:8081 - - connect: - image: cnfldemos/cp-server-connect-datagen:0.6.2-7.5.0 - hostname: connect - container_name: connect - depends_on: - - broker - - schema-registry - ports: - - "8083:8083" - environment: - CONNECT_BOOTSTRAP_SERVERS: 'broker:29092' - CONNECT_REST_ADVERTISED_HOST_NAME: connect - CONNECT_GROUP_ID: compose-connect-group - CONNECT_CONFIG_STORAGE_TOPIC: docker-connect-configs - CONNECT_CONFIG_STORAGE_REPLICATION_FACTOR: 1 - CONNECT_OFFSET_FLUSH_INTERVAL_MS: 10000 - CONNECT_OFFSET_STORAGE_TOPIC: docker-connect-offsets - CONNECT_OFFSET_STORAGE_REPLICATION_FACTOR: 1 - CONNECT_STATUS_STORAGE_TOPIC: docker-connect-status - CONNECT_STATUS_STORAGE_REPLICATION_FACTOR: 1 - CONNECT_KEY_CONVERTER: org.apache.kafka.connect.storage.StringConverter - CONNECT_VALUE_CONVERTER: io.confluent.connect.avro.AvroConverter - CONNECT_VALUE_CONVERTER_SCHEMA_REGISTRY_URL: http://schema-registry:8081 - # CLASSPATH required due to CC-2422 - CLASSPATH: /usr/share/java/monitoring-interceptors/monitoring-interceptors-7.5.3.jar - CONNECT_PRODUCER_INTERCEPTOR_CLASSES: "io.confluent.monitoring.clients.interceptor.MonitoringProducerInterceptor" - CONNECT_CONSUMER_INTERCEPTOR_CLASSES: "io.confluent.monitoring.clients.interceptor.MonitoringConsumerInterceptor" - CONNECT_PLUGIN_PATH: "/usr/share/java,/usr/share/confluent-hub-components" - CONNECT_LOG4J_LOGGERS: org.apache.zookeeper=ERROR,org.I0Itec.zkclient=ERROR,org.reflections=ERROR - - control-center: - image: confluentinc/cp-enterprise-control-center:7.5.3 - hostname: control-center - container_name: control-center - depends_on: - - broker - - schema-registry - - connect - - ksqldb-server - ports: - - "9021:9021" - environment: - CONTROL_CENTER_BOOTSTRAP_SERVERS: 'broker:29092' - CONTROL_CENTER_CONNECT_CONNECT-DEFAULT_CLUSTER: 'connect:8083' - CONTROL_CENTER_KSQL_KSQLDB1_URL: "http://ksqldb-server:8088" - CONTROL_CENTER_KSQL_KSQLDB1_ADVERTISED_URL: "http://localhost:8088" - CONTROL_CENTER_SCHEMA_REGISTRY_URL: "http://schema-registry:8081" - CONTROL_CENTER_REPLICATION_FACTOR: 1 - CONTROL_CENTER_INTERNAL_TOPICS_PARTITIONS: 1 - CONTROL_CENTER_MONITORING_INTERCEPTOR_TOPIC_PARTITIONS: 1 - CONFLUENT_METRICS_TOPIC_REPLICATION: 1 - PORT: 9021 - - ksqldb-server: - image: confluentinc/cp-ksqldb-server:7.5.3 - hostname: ksqldb-server - container_name: ksqldb-server - depends_on: - - broker - - connect - ports: - - "8088:8088" - environment: - KSQL_CONFIG_DIR: "/etc/ksql" - KSQL_BOOTSTRAP_SERVERS: "broker:29092" - KSQL_HOST_NAME: ksqldb-server - KSQL_LISTENERS: "http://0.0.0.0:8088" - KSQL_CACHE_MAX_BYTES_BUFFERING: 0 - KSQL_KSQL_SCHEMA_REGISTRY_URL: "http://schema-registry:8081" - KSQL_PRODUCER_INTERCEPTOR_CLASSES: "io.confluent.monitoring.clients.interceptor.MonitoringProducerInterceptor" - KSQL_CONSUMER_INTERCEPTOR_CLASSES: "io.confluent.monitoring.clients.interceptor.MonitoringConsumerInterceptor" - KSQL_KSQL_CONNECT_URL: "http://connect:8083" - KSQL_KSQL_LOGGING_PROCESSING_TOPIC_REPLICATION_FACTOR: 1 - KSQL_KSQL_LOGGING_PROCESSING_TOPIC_AUTO_CREATE: 'true' - KSQL_KSQL_LOGGING_PROCESSING_STREAM_AUTO_CREATE: 'true' - - # ksqldb-cli: - # image: confluentinc/cp-ksqldb-cli:7.5.3 - # container_name: ksqldb-cli - # depends_on: - # - broker - # - connect - # - ksqldb-server - # entrypoint: /bin/sh - # tty: true - - # ksql-datagen: - # image: confluentinc/ksqldb-examples:7.5.3 - # hostname: ksql-datagen - # container_name: ksql-datagen - # depends_on: - # - ksqldb-server - # - broker - # - schema-registry - # - connect - # command: "bash -c 'echo Waiting for Kafka to be ready... && \ - # cub kafka-ready -b broker:29092 1 40 && \ - # echo Waiting for Confluent Schema Registry to be ready... && \ - # cub sr-ready schema-registry 8081 40 && \ - # echo Waiting a few seconds for topic creation to finish... && \ - # sleep 11 && \ - # tail -f /dev/null'" - # environment: - # KSQL_CONFIG_DIR: "/etc/ksql" - # STREAMS_BOOTSTRAP_SERVERS: broker:29092 - # STREAMS_SCHEMA_REGISTRY_HOST: schema-registry - # STREAMS_SCHEMA_REGISTRY_PORT: 8081 - - rest-proxy: - image: confluentinc/cp-kafka-rest:7.5.3 + KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 + # Consumer groups are created per-run, so skip the 3s default rebalance debounce + KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0 + # Tests create their topics explicitly; test_invalid_topic depends on this being off + KAFKA_AUTO_CREATE_TOPICS_ENABLE: "false" + healthcheck: + test: ["CMD-SHELL", "kafka-broker-api-versions --bootstrap-server localhost:9092"] + interval: 5s + timeout: 10s + retries: 24 depends_on: - - broker - - schema-registry - ports: - - 8082:8082 - hostname: rest-proxy - container_name: rest-proxy - environment: - KAFKA_REST_HOST_NAME: rest-proxy - KAFKA_REST_BOOTSTRAP_SERVERS: 'broker:29092' - KAFKA_REST_LISTENERS: "http://0.0.0.0:8082" - KAFKA_REST_SCHEMA_REGISTRY_URL: 'http://schema-registry:8081' \ No newline at end of file + zoo1: + condition: service_healthy diff --git a/conda/dev-environment-unix.yml b/conda/dev-environment-unix.yml index fadb0b141..9ff9f08e8 100644 --- a/conda/dev-environment-unix.yml +++ b/conda/dev-environment-unix.yml @@ -17,7 +17,6 @@ dependencies: - flex - graphviz - gtest - - httpx>=0.20,<1 - libarrow<24 - libboost>=1.80.0 - libboost-headers>=1.80.0 @@ -41,6 +40,7 @@ dependencies: - pytest-sugar - python>=3.10,<3.15 - python-build + - python-confluent-kafka - python-graphviz - python-rapidjson - pytz diff --git a/conda/dev-environment-win.yml b/conda/dev-environment-win.yml index 65ca65cb4..4ac9e69e4 100644 --- a/conda/dev-environment-win.yml +++ b/conda/dev-environment-win.yml @@ -17,7 +17,6 @@ dependencies: # - flex # not available on windows - graphviz - gtest - - httpx>=0.20,<1 - libarrow<24 - libboost>=1.80.0 - libboost-headers>=1.80.0 @@ -41,6 +40,7 @@ dependencies: - pytest-sugar - python>=3.10,<3.14 - python-build + - python-confluent-kafka - python-graphviz - python-rapidjson - pytz diff --git a/cpp/csp/engine/Struct.h b/cpp/csp/engine/Struct.h index 7082fefc9..f4c4f5958 100644 --- a/cpp/csp/engine/Struct.h +++ b/cpp/csp/engine/Struct.h @@ -737,7 +737,7 @@ std::shared_ptr::type> StructMeta::getMetaField( std::shared_ptr::type> typedfield = std::dynamic_pointer_cast::type>( field_ ); if( !typedfield ) CSP_THROW( TypeError, expectedtype << " - provided struct type " << name() << " expected type " << CspType::Type::fromCType::type << " for field " << fieldname - << " but got type " << field_ -> type() -> type() << " for " << expectedtype ); + << " but got type " << field_ -> type() -> type() ); return typedfield; } diff --git a/csp/adapters/kafka.py b/csp/adapters/kafka.py index daa3282d0..45032dbc8 100644 --- a/csp/adapters/kafka.py +++ b/csp/adapters/kafka.py @@ -20,7 +20,18 @@ from csp.impl.wiring import ReplayMode, input_adapter_def, output_adapter_def, status_adapter_def from csp.lib import _kafkaadapterimpl -_ = BytesMessageProtoMapper, DateTimeType, JSONTextMessageMapper, RawBytesMessageMapper, RawTextMessageMapper +__all__ = ( + "KafkaAdapterManager", + "KafkaStartOffset", + "KafkaStatusMessageType", + # re-exported from csp.adapters.utils for backwards compatibility + "BytesMessageProtoMapper", + "DateTimeType", + "JSONTextMessageMapper", + "RawBytesMessageMapper", + "RawTextMessageMapper", +) + T = TypeVar("T") diff --git a/csp/tests/adapters/conftest.py b/csp/tests/adapters/conftest.py index 62046fa4a..280cb3603 100644 --- a/csp/tests/adapters/conftest.py +++ b/csp/tests/adapters/conftest.py @@ -1,3 +1,5 @@ +from uuid import uuid4 + import pytest from csp.adapters.kafka import KafkaAdapterManager @@ -5,15 +7,16 @@ @pytest.fixture(scope="module", autouse=True) def kafkabroker(): + # Defined in ci/kafka/docker-compose.yml return "localhost:9092" @pytest.fixture(scope="module", autouse=True) def kafkaadapterkwargs(kafkabroker): - return dict(broker=kafkabroker, group_id="group.id123", rd_kafka_conf_options={"allow.auto.create.topics": "true"}) + # Unique group id so a rerun never inherits committed offsets from a previous run + return dict(broker=kafkabroker, group_id=f"csp.test.{uuid4()}") @pytest.fixture(scope="module", autouse=True) def kafkaadapter(kafkaadapterkwargs): - _kafkaadapter = KafkaAdapterManager(**kafkaadapterkwargs) - return _kafkaadapter + return KafkaAdapterManager(**kafkaadapterkwargs) diff --git a/csp/tests/adapters/kafka_utils.py b/csp/tests/adapters/kafka_utils.py index 98544cd34..42ccbd3e1 100644 --- a/csp/tests/adapters/kafka_utils.py +++ b/csp/tests/adapters/kafka_utils.py @@ -1,13 +1,16 @@ -import httpx - - -def _precreate_topic(topic): - """Since we test against confluent kafka, just use the kafka rest addon""" - rest_broker = "http://localhost:8082" - cluster_info = httpx.get(f"{rest_broker}/v3/clusters") - cluster_id = cluster_info.json()["data"][0]["cluster_id"] - resp = httpx.post(f"{rest_broker}/v3/clusters/{cluster_id}/topics", json={"topic_name": topic}) - if resp.status_code != 201 and "already exists" not in resp.content.decode("utf8"): - raise Exception( - f"Could not create topic {topic} on cluster {rest_broker}/v3/clusters/{cluster_id}/topics - received {resp.content}" - ) +__all__ = ("create_topic",) + + +def create_topic(broker, topic): + """Create `topic` and block until the broker acknowledges it. + + Creation is done out of band rather than by publishing a warm-up message, so that no test data + lands on the topic and so tests do not depend on broker-side auto-creation (disabled in + ci/kafka/docker-compose.yml so test_invalid_topic can exercise the failure path). + """ + # Imported lazily so collection does not require confluent-kafka when the kafka tests are skipped + from confluent_kafka.admin import AdminClient, NewTopic + + admin = AdminClient({"bootstrap.servers": broker}) + for _, future in admin.create_topics([NewTopic(topic, num_partitions=1, replication_factor=1)]).items(): + future.result() diff --git a/csp/tests/adapters/test_kafka.py b/csp/tests/adapters/test_kafka.py index 6cc15388f..dd280bf92 100644 --- a/csp/tests/adapters/test_kafka.py +++ b/csp/tests/adapters/test_kafka.py @@ -6,17 +6,11 @@ import csp from csp import ts -from csp.adapters.kafka import ( - DateTimeType, - JSONTextMessageMapper, - KafkaAdapterManager, - KafkaStartOffset, - RawBytesMessageMapper, - RawTextMessageMapper, -) +from csp.adapters.kafka import KafkaAdapterManager, KafkaStartOffset +from csp.adapters.utils import DateTimeType, JSONTextMessageMapper, RawBytesMessageMapper, RawTextMessageMapper from csp.utils.datetime import utc_now -from .kafka_utils import _precreate_topic +from .kafka_utils import create_topic class MyData(csp.Struct): @@ -69,7 +63,10 @@ class MetaSubData(csp.Struct): class TestKafka: @pytest.mark.skipif(not os.environ.get("CSP_TEST_KAFKA"), reason="Skipping kafka adapter tests") - def test_metadata(self, kafkaadapter): + def test_metadata(self, kafkaadapter, kafkabroker): + topic = f"test.metadata.{os.getpid()}" + create_topic(kafkabroker, topic) + def graph(count: int): msg_mapper = JSONTextMessageMapper(datetime_type=DateTimeType.UINT64_MICROS) @@ -84,19 +81,17 @@ def graph(count: int): "timestamp": "mapped_timestamp", } - topic = f"test.metadata.{os.getpid()}" - _precreate_topic(topic) subKey = "foo" pubKey = ["mapped_a", "mapped_b", "mapped_c"] - c = csp.count(csp.timer(timedelta(seconds=0.1))) + # Publish slowly enough that the consumer reaches end-of-partition and flips to live + c = csp.count(csp.timer(timedelta(seconds=1))) t = csp.sample(c, csp.const("foo")) pubStruct = MetaPubData.collectts( mapped_a=MetaSubStruct.collectts(mapped_b=MetaTextStruct.collectts(mapped_c=t)), mapped_count=c ) - # csp.print('pub', pubStruct) kafkaadapter.publish(msg_mapper, topic, pubKey, pubStruct, field_map=pub_field_map) sub_data = kafkaadapter.subscribe( @@ -110,7 +105,7 @@ def graph(count: int): ) csp.add_graph_output("sub_data", sub_data) - # csp.print('sub', sub_data) + # Wait for at least count ticks and until we get a live tick done_flag = csp.count(sub_data) >= count done_flag = csp.and_(done_flag, sub_data.mapped_live == True) # noqa: E712 @@ -119,17 +114,18 @@ def graph(count: int): count = 5 results = csp.run(graph, count, starttime=utc_now(), endtime=timedelta(seconds=30), realtime=True) - assert len(results["sub_data"]) >= 5 - print(results) + assert len(results["sub_data"]) >= count + for result in results["sub_data"]: assert result[1].mapped_partition >= 0 assert result[1].mapped_offset >= 0 assert result[1].mapped_live is not None assert result[1].mapped_timestamp < utc_now() + # last record should be live (first may or may not be live depending on timing) assert results["sub_data"][-1][1].mapped_live @pytest.mark.skipif(not os.environ.get("CSP_TEST_KAFKA"), reason="Skipping kafka adapter tests") - def test_basic(self, kafkaadapter): + def test_basic(self, kafkaadapter, kafkabroker): @csp.node def curtime(x: ts[object]) -> ts[datetime]: if csp.ticked(x): @@ -140,6 +136,7 @@ def graph(symbols: list, count: int): csp.timer(timedelta(seconds=0.2), True), csp.delay(csp.timer(timedelta(seconds=0.2), False), timedelta(seconds=0.1)), ) + i = csp.count(csp.timer(timedelta(seconds=0.15))) d = csp.count(csp.timer(timedelta(seconds=0.2))) / 2.0 s = csp.sample(csp.timer(timedelta(seconds=0.4)), csp.const("STRING")) @@ -153,8 +150,7 @@ def graph(symbols: list, count: int): struct_field_map = {"b": "b2", "i": "i2", "d": "d2", "s": "s2", "dt": "dt2", "date": "date2"} done_flags = [] - topic = f"mktdata.{os.getpid()}" - _precreate_topic(topic) + for symbol in symbols: kafkaadapter.publish(msg_mapper, topic, symbol, b, field_map="b") kafkaadapter.publish(msg_mapper, topic, symbol, i, field_map="i") @@ -181,8 +177,6 @@ def graph(symbols: list, count: int): ) csp.add_graph_output(f"pall_{symbol}", pub_data) - # csp.print('status', kafkaadapter.status()) - sub_data = kafkaadapter.subscribe( ts_type=SubData, msg_mapper=msg_mapper, @@ -190,9 +184,6 @@ def graph(symbols: list, count: int): key=symbol, push_mode=csp.PushMode.NON_COLLAPSING, ) - - sub_data = csp.firstN(sub_data, count) - csp.add_graph_output(f"sall_{symbol}", sub_data) done_flag = csp.count(sub_data) == count @@ -203,20 +194,25 @@ def graph(symbols: list, count: int): stop = csp.filter(stop, stop) csp.stop_engine(stop) + topic = f"mktdata.{os.getpid()}" + create_topic(kafkabroker, topic) symbols = ["AAPL", "MSFT"] - count = 100 + count = 50 results = csp.run(graph, symbols, count, starttime=utc_now(), endtime=timedelta(seconds=30), realtime=True) for symbol in symbols: - pub = results[f"pall_{symbol}"] - sub = results[f"sall_{symbol}"] + pub = [v[1] for v in results[f"pall_{symbol}"]] + sub = [v[1] for v in results[f"sall_{symbol}"]] + # The subscriber may miss a prefix of the stream while its consumer group is being + # assigned, so align on the first record it did see rather than assuming it saw all. assert len(sub) == count - assert [v[1] for v in sub] == [v[1] for v in pub[:count]] + start = pub.index(sub[0]) + assert sub == pub[start : start + count] @pytest.mark.skipif(not os.environ.get("CSP_TEST_KAFKA"), reason="Skipping kafka adapter tests") def test_start_offsets(self, kafkaadapter, kafkabroker): topic = f"test_start_offsets.{os.getpid()}" - _precreate_topic(topic) + create_topic(kafkabroker, topic) msg_mapper = JSONTextMessageMapper(datetime_type=DateTimeType.UINT64_MICROS) count = 10 @@ -228,7 +224,6 @@ def pub_graph(): stop = csp.count(struct) == count stop = csp.filter(stop, stop) csp.stop_engine(stop) - # csp.print('pub', struct) csp.run(pub_graph, starttime=utc_now(), endtime=timedelta(seconds=30), realtime=True) @@ -247,9 +242,6 @@ def get_times_graph(): csp.stop_engine(csp.filter(stop, stop)) csp.add_graph_output("data", data) - # csp.print('sub', data) - # csp.print('status', kafkaadapter.status()) - all_data = csp.run(get_times_graph, starttime=utc_now(), endtime=timedelta(seconds=30), realtime=True)["data"] min_time = all_data[0][1].dt @@ -267,8 +259,6 @@ def get_data(start_offset, expected_count): csp.stop_engine(csp.filter(stop, stop)) csp.add_graph_output("data", data) - # csp.print('data', data) - res = csp.run( get_data, KafkaStartOffset.EARLIEST, @@ -277,7 +267,6 @@ def get_data(start_offset, expected_count): endtime=timedelta(seconds=30), realtime=True, )["data"] - # print(res) # If we playback from earliest but start "now", all data should still arrive but as realtime ticks assert len(res) == 10 @@ -292,7 +281,7 @@ def get_data(start_offset, expected_count): assert len(res) == 0 res = csp.run( - get_data, KafkaStartOffset.START_TIME, 10, starttime=min_time, endtime=timedelta(seconds=30), realtime=True + get_data, KafkaStartOffset.START_TIME, 10, starttime=min_time, endtime=timedelta(seconds=10), realtime=True )["data"] assert len(res) == 10 @@ -308,12 +297,12 @@ def get_data(start_offset, expected_count): assert len(res) == len(expected) res = csp.run( - get_data, timedelta(seconds=0), len(expected), starttime=stime, endtime=timedelta(seconds=30), realtime=True + get_data, timedelta(seconds=0), len(expected), starttime=stime, endtime=timedelta(seconds=10), realtime=True )["data"] assert len(res) == len(expected) @pytest.mark.skipif(not os.environ.get("CSP_TEST_KAFKA"), reason="Skipping kafka adapter tests") - def test_raw_pubsub(self, kafkaadapter): + def test_raw_pubsub(self, kafkaadapter, kafkabroker): @csp.node def data(x: ts[object]) -> ts[bytes]: if csp.ticked(x): @@ -329,15 +318,10 @@ def graph(symbols: list, count: int): msg_mapper = RawBytesMessageMapper() done_flags = [] - topic = f"test_str.{os.getpid()}" - _precreate_topic(topic) for symbol in symbols: - topic = f"test_str.{os.getpid()}" kafkaadapter.publish(msg_mapper, topic, symbol, d) csp.add_graph_output(f"pub_{symbol}", d) - # csp.print('status', kafkaadapter.status()) - sub_data = kafkaadapter.subscribe( ts_type=SubData, msg_mapper=RawTextMessageMapper(), @@ -356,14 +340,12 @@ def graph(symbols: list, count: int): push_mode=csp.PushMode.NON_COLLAPSING, ) - sub_data = csp.firstN(sub_data.msg, count) - sub_data_bytes = csp.firstN(sub_data_bytes, count) - - # csp.print('sub', sub_data) - csp.add_graph_output(f"sub_{symbol}", sub_data) + csp.add_graph_output(f"sub_{symbol}", sub_data.msg) csp.add_graph_output(f"sub_bytes_{symbol}", sub_data_bytes) - done_flag = csp.count(sub_data) + csp.count(sub_data_bytes) == count * 2 + # Wait for count messages on both subscribers + done_flag = csp.count(sub_data) >= count + done_flag = csp.and_(done_flag, csp.count(sub_data_bytes) >= count) done_flag = csp.filter(done_flag, done_flag) done_flags.append(done_flag) @@ -371,29 +353,35 @@ def graph(symbols: list, count: int): stop = csp.filter(stop, stop) csp.stop_engine(stop) + topic = f"test_str.{os.getpid()}" + create_topic(kafkabroker, topic) + symbols = ["AAPL", "MSFT"] count = 10 results = csp.run(graph, symbols, count, starttime=utc_now(), endtime=timedelta(seconds=30), realtime=True) - # print(results) for symbol in symbols: - pub = results[f"pub_{symbol}"] - sub = results[f"sub_{symbol}"] - sub_bytes = results[f"sub_bytes_{symbol}"] - - assert len(sub) == count - assert [v[1] for v in sub] == [v[1] for v in pub[:count]] - assert [v[1] for v in sub_bytes] == [v[1] for v in pub[:count]] + pub = [v[1] for v in results[f"pub_{symbol}"]] + sub = [v[1] for v in results[f"sub_{symbol}"]] + sub_bytes = [v[1] for v in results[f"sub_bytes_{symbol}"]] + + # Align on the first record each subscriber saw, then require an exact contiguous run. + # Comparing sequences (rather than set membership) is what catches loss, duplication + # and reordering. + for received in (sub, sub_bytes): + assert len(received) >= count + start = pub.index(received[0]) + assert received == pub[start : start + len(received)] @pytest.mark.skipif(not os.environ.get("CSP_TEST_KAFKA"), reason="Skipping kafka adapter tests") def test_invalid_topic(self, kafkaadapterkwargs): class SubData(csp.Struct): msg: str + # Relies on broker-side auto.create.topics.enable=false (ci/kafka/docker-compose.yml) kafkaadapter1 = KafkaAdapterManager(**kafkaadapterkwargs) # Was a bug where engine would stall def graph_sub(): - # csp.print('status', kafkaadapter.status()) return kafkaadapter1.subscribe( ts_type=SubData, msg_mapper=RawTextMessageMapper(), field_map={"": "msg"}, topic="foobar", key="none" ) @@ -401,6 +389,7 @@ def graph_sub(): # With bug this would deadlock with pytest.raises(RuntimeError): csp.run(graph_sub, starttime=utc_now(), endtime=timedelta(seconds=2), realtime=True) + kafkaadapter2 = KafkaAdapterManager(**kafkaadapterkwargs) def graph_pub(): @@ -442,15 +431,13 @@ def graph_pub(): csp.run(graph_pub, starttime=utc_now(), endtime=timedelta(seconds=2), realtime=True) @pytest.mark.skipif(not os.environ.get("CSP_TEST_KAFKA"), reason="Skipping kafka adapter tests") - def test_meta_field_map_tick_timestamp_from_field(self, kafkaadapterkwargs): + def test_meta_field_map_tick_timestamp_from_field(self, kafkaadapter): class SubData(csp.Struct): msg: str dt: datetime - kafkaadapter1 = KafkaAdapterManager(**kafkaadapterkwargs) - def graph_sub(): - return kafkaadapter1.subscribe( + return kafkaadapter.subscribe( ts_type=SubData, msg_mapper=RawTextMessageMapper(), meta_field_map={"timestamp": "dt"}, @@ -482,7 +469,7 @@ class BasicData(csp.Struct): b: bool topic = f"test_burst.{os.getpid()}" - _precreate_topic(topic) + create_topic(kafkabroker, topic) msg_mapper = JSONTextMessageMapper(datetime_type=DateTimeType.UINT64_MICROS) count = 10 diff --git a/csp/tests/adapters/test_status.py b/csp/tests/adapters/test_status.py index 9ad5e7147..0977478cb 100644 --- a/csp/tests/adapters/test_status.py +++ b/csp/tests/adapters/test_status.py @@ -1,24 +1,25 @@ import os -from datetime import datetime, timedelta +from datetime import timedelta import pytest import csp from csp import ts -from csp.adapters.kafka import DateTimeType, JSONTextMessageMapper, KafkaStatusMessageType +from csp.adapters.kafka import KafkaStatusMessageType from csp.adapters.status import Level +from csp.adapters.utils import DateTimeType, JSONTextMessageMapper from csp.utils.datetime import utc_now -from .kafka_utils import _precreate_topic +from .kafka_utils import create_topic class SubData(csp.Struct): a: bool -class TestStatus: +class TestStatusKafka: @pytest.mark.skipif(not os.environ.get("CSP_TEST_KAFKA"), reason="Skipping kafka adapter tests") - def test_basic(self, kafkaadapter): + def test_basic(self, kafkaadapter, kafkabroker): topic = f"csp.unittest.{os.getpid()}" key = "test_status" @@ -43,7 +44,7 @@ def graph(): done_flag = csp.count(status) == 1 csp.stop_engine(done_flag) - _precreate_topic(topic) + create_topic(kafkabroker, topic) results = csp.run(graph, starttime=utc_now(), endtime=timedelta(seconds=10), realtime=True) status = results["status"][0][1] assert status.status_code == KafkaStatusMessageType.MSG_RECV_ERROR diff --git a/examples/03_using_adapters/kafka/e1_kafka.py b/examples/03_using_adapters/kafka/e1_kafka.py index a19ace9ca..4ac1eb2e9 100644 --- a/examples/03_using_adapters/kafka/e1_kafka.py +++ b/examples/03_using_adapters/kafka/e1_kafka.py @@ -4,10 +4,10 @@ import csp from csp import ts from csp.adapters.kafka import ( + BytesMessageProtoMapper, DateTimeType, JSONTextMessageMapper, KafkaAdapterManager, - ProtoMessageMapper, RawTextMessageMapper, ) from csp.utils.datetime import utc_now @@ -142,7 +142,7 @@ def proto_graph(): "px": "price", } - msg_mapper = ProtoMessageMapper( + msg_mapper = BytesMessageProtoMapper( proto_directory="/tmp", proto_filename="fxspotstream.proto", proto_message="Snapshot" ) @@ -157,7 +157,7 @@ def proto_graph_multiple_subscribers(): topic = "test2" - msg_mapper = ProtoMessageMapper( + msg_mapper = BytesMessageProtoMapper( proto_directory="/tmp", proto_filename="fxspotstream.proto", proto_message="Snapshot" ) diff --git a/pyproject.toml b/pyproject.toml index 96396ee30..de892a3ed 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -81,9 +81,10 @@ develop = [ "graphviz", "pillow", # adapters - "httpx>=0.20,<1", # kafka + "confluent-kafka", # kafka test topic setup "perspective-python>=2", # perspective - "anywidget", # perspective + "anywidget", # perspective (perspective-python 3.x) + "ipywidgets", # perspective (csp.impl.pandas_perspective) "polars", # parquet "psutil", # test_engine/test_history "sqlalchemy", # db @@ -96,21 +97,25 @@ showgraph = [ "pillow", ] test = [ - "graphviz", - "pillow", "pytest", "pytest-asyncio", "pytest-cov", "pytest-sugar", - "httpx>=0.20,<1", - "perspective-python", - "polars", - "psutil", - "requests", - "slack-sdk>=3", - "sqlalchemy", - "threadpoolctl", - "tornado", + # showgraph + "graphviz", + "pillow", + # adapters + "confluent-kafka", # kafka test topic setup + "perspective-python>=2", # perspective + "anywidget", # perspective (perspective-python 3.x) + "ipywidgets", # perspective (csp.impl.pandas_perspective) + "polars", # parquet + "psutil", # test_engine/test_history + "sqlalchemy", # db + "threadpoolctl", # test_random + "requests", # profiler + "tornado", # profiler, perspective, websocket + "python-rapidjson", # websocket ] symphony = [ "csp-adapter-symphony", From be0928b5fba0082f76fdaf457143b0f0ae8b56fa Mon Sep 17 00:00:00 2001 From: Tim Paine <3105306+timkpaine@users.noreply.github.com> Date: Fri, 14 Aug 2026 18:43:26 -0400 Subject: [PATCH 2/2] Make Kafka adapter tolerate transient outages and report missing topics The tests bootstrapped from "localhost", which also resolves to ::1 where the container publishes nothing, so librdkafka failed whichever address it picked next. Point them at 127.0.0.1, matching the published port and the advertised listener. A single all-brokers-down report then killed the engine, though librdkafka reports it on every failed connection round and reconnects on its own. Escalation now waits for the brokers to stay down for broker_down_tolerance, checked on the poll threads rather than on the report, since librdkafka may not report again for seconds. Shutdown no longer waits on a flush that cannot complete. A topic that the broker will not auto-create left a subscriber looking idle forever. The consumer reports that through poll rather than the event callback, so escalate it there, and on the publisher through the delivery report. Signed-off-by: Tim Paine <3105306+timkpaine@users.noreply.github.com> --- .../adapters/kafka/KafkaAdapterManager.cpp | 59 +++++++++++++++++-- cpp/csp/adapters/kafka/KafkaAdapterManager.h | 12 ++++ cpp/csp/adapters/kafka/KafkaConsumer.cpp | 14 +++++ csp/adapters/kafka.py | 4 ++ csp/tests/adapters/conftest.py | 5 +- csp/tests/adapters/test_kafka.py | 9 ++- 6 files changed, 94 insertions(+), 9 deletions(-) diff --git a/cpp/csp/adapters/kafka/KafkaAdapterManager.cpp b/cpp/csp/adapters/kafka/KafkaAdapterManager.cpp index b6251364d..1bae12fa9 100644 --- a/cpp/csp/adapters/kafka/KafkaAdapterManager.cpp +++ b/cpp/csp/adapters/kafka/KafkaAdapterManager.cpp @@ -23,6 +23,14 @@ INIT_CSP_ENUM( csp::adapters::kafka::KafkaStatusMessageType, namespace csp::adapters::kafka { +//A topic that does not exist is only reported this way when the broker will not auto-create it. +//The producer reports it after topic.metadata.propagation.max.ms rather than immediately. +static bool isUnknownTopic( int err ) +{ + return err == RdKafka::ErrorCode::ERR_UNKNOWN_TOPIC_OR_PART || + err == RdKafka::ErrorCode::ERR__UNKNOWN_TOPIC; +} + class DeliveryReportCb : public RdKafka::DeliveryReportCb { public: @@ -38,7 +46,12 @@ class DeliveryReportCb : public RdKafka::DeliveryReportCb { std::string msg = "KafkaPublisher: Message delivery failed for topic " + message.topic_name() + ". Failure: " + message.errstr(); m_adapterManager -> pushStatus( StatusLevel::ERROR, KafkaStatusMessageType::MSG_DELIVERY_FAILED, msg ); + + if( isUnknownTopic( message.err() ) ) + m_adapterManager -> forceShutdown( msg ); } + else + m_adapterManager -> onBrokerActivity(); } private: KafkaAdapterManager * m_adapterManager; @@ -60,13 +73,17 @@ class EventCb : public RdKafka::EventCb if( event.type() == RdKafka::Event::EVENT_ERROR ) { //We shutdown the app if its a fatal error OR if its an authentication issue which has plagued users multiple times - //Adding ERR__ALL_BROKERS_DOWN which happens when all brokers are down if( event.fatal() || - event.err() == RdKafka::ErrorCode::ERR__AUTHENTICATION || - event.err() == RdKafka::ErrorCode::ERR__ALL_BROKERS_DOWN ) + event.err() == RdKafka::ErrorCode::ERR__AUTHENTICATION ) { m_adapterManager -> forceShutdown( RdKafka::err2str( ( RdKafka::ErrorCode ) event.err() ) + event.str() ); } + //librdkafka reports this on every failed connection round and reconnects on its own, so a + //single report says nothing about whether the brokers are really gone + else if( event.err() == RdKafka::ErrorCode::ERR__ALL_BROKERS_DOWN ) + { + m_adapterManager -> onBrokersDown( RdKafka::err2str( ( RdKafka::ErrorCode ) event.err() ) + event.str() ); + } } } @@ -77,10 +94,12 @@ class EventCb : public RdKafka::EventCb KafkaAdapterManager::KafkaAdapterManager( csp::Engine * engine, const Dictionary & properties ) : AdapterManager( engine ), m_consumerIdx( 0 ), m_producerPollThreadActive( false ), - m_unrecoverableError( false ) + m_unrecoverableError( false ), + m_brokersDownSince( 0 ) { m_maxThreads = properties.get( "max_threads" ); m_pollTimeoutMs = properties.get( "poll_timeout" ).asMilliseconds(); + m_brokersDownTolerance = properties.get( "broker_down_tolerance" ).asNanoseconds(); m_eventCb = std::make_unique( this ); m_producerCb = std::make_unique( this ); @@ -133,6 +152,33 @@ void KafkaAdapterManager::setConfProperties( RdKafka::Conf * conf, const Diction } } +void KafkaAdapterManager::onBrokersDown( const std::string & err ) +{ + int64_t since = 0; + + //Only the first report starts the clock. Escalation is left to checkBrokersDown() on the poll + //threads rather than done here, since librdkafka may not report again for seconds. + if( m_brokersDownSince.compare_exchange_strong( since, DateTime::now().asNanoseconds() ) ) + { + std::lock_guard guard( m_brokersDownLock ); + m_brokersDownError = err; + } +} + +void KafkaAdapterManager::checkBrokersDown() +{ + int64_t since = m_brokersDownSince.load( std::memory_order_relaxed ); + if( since == 0 || DateTime::now().asNanoseconds() - since < m_brokersDownTolerance ) + return; + + std::string err; + { + std::lock_guard guard( m_brokersDownLock ); + err = m_brokersDownError; + } + forceShutdown( err ); +} + void KafkaAdapterManager::forceShutdown( const std::string & err ) { m_unrecoverableError = true; // So we can alert the producer to stop trying to flush @@ -228,6 +274,7 @@ void KafkaAdapterManager::pollProducers() while( m_producerPollThreadActive ) { m_producer -> poll( m_pollTimeoutMs ); + checkBrokersDown(); } try @@ -235,7 +282,9 @@ void KafkaAdapterManager::pollProducers() while( true ) { auto rc = m_producer -> flush( 5000 ); - if( !rc || m_unrecoverableError ) + //Waiting on a flush that cannot complete is what used to hang shutdown, so give up as + //soon as there is nothing to flush to + if( !rc || m_unrecoverableError || brokersDown() ) break; if( rc != RdKafka::ERR__TIMED_OUT ) diff --git a/cpp/csp/adapters/kafka/KafkaAdapterManager.h b/cpp/csp/adapters/kafka/KafkaAdapterManager.h index c095d026a..fbb275125 100644 --- a/cpp/csp/adapters/kafka/KafkaAdapterManager.h +++ b/cpp/csp/adapters/kafka/KafkaAdapterManager.h @@ -8,6 +8,7 @@ #include #include #include +#include #include #include #include @@ -75,6 +76,11 @@ class KafkaAdapterManager final : public csp::AdapterManager void forceShutdown( const std::string & err ); + void onBrokersDown( const std::string & err ); + void checkBrokersDown(); + bool brokersDown() const { return m_brokersDownSince.load( std::memory_order_relaxed ) != 0; } + void onBrokerActivity() { m_brokersDownSince.store( 0, std::memory_order_relaxed ); } + void markConsumerReplayDone( KafkaConsumer * consumer, const std::string & topic ); void onMessage( RdKafka::Message * msg ) const; @@ -132,6 +138,12 @@ class KafkaAdapterManager final : public csp::AdapterManager std::atomic m_producerPollThreadActive; std::atomic m_unrecoverableError; + //Nanoseconds since the brokers were first reported down, 0 while they are believed up + std::atomic m_brokersDownSince; + int64_t m_brokersDownTolerance; + std::mutex m_brokersDownLock; + std::string m_brokersDownError; + std::unique_ptr m_consumerConf; std::unique_ptr m_producerConf; Dictionary::Value m_startOffsetProperty; diff --git a/cpp/csp/adapters/kafka/KafkaConsumer.cpp b/cpp/csp/adapters/kafka/KafkaConsumer.cpp index de353e639..2031d12cf 100644 --- a/cpp/csp/adapters/kafka/KafkaConsumer.cpp +++ b/cpp/csp/adapters/kafka/KafkaConsumer.cpp @@ -196,6 +196,8 @@ void KafkaConsumer::poll() { std::unique_ptr msg( m_consumer -> consume( m_mgr -> pollTimeoutMs() ) ); + m_mgr -> checkBrokersDown(); + if( msg -> err() == RdKafka::ERR__TIMED_OUT ) continue; @@ -208,8 +210,20 @@ void KafkaConsumer::poll() continue; } + //Only reported when the broker will not auto-create the topic. Waiting for a topic + //that will never appear looks identical to an idle subscription, so surface it. + if( msg -> err() == RdKafka::ERR_UNKNOWN_TOPIC_OR_PART ) [[unlikely]] + { + m_mgr -> forceShutdown( RdKafka::err2str( msg -> err() ) + " error: " + msg -> errstr() ); + continue; + } + if( msg -> err() == RdKafka::ERR_NO_ERROR && msg -> len() ) + { + //Proof the brokers came back, whatever was reported down before + m_mgr -> onBrokerActivity(); m_mgr -> onMessage( msg.get() ); + } //Not sure why, but it looks like we repeatedly get EOF callbacks even after the original one //may want to look into this. Not an issue in practice, but seems like unnecessary overhead else if( msg -> err() == RdKafka::ERR__PARTITION_EOF ) diff --git a/csp/adapters/kafka.py b/csp/adapters/kafka.py index 45032dbc8..0dc368838 100644 --- a/csp/adapters/kafka.py +++ b/csp/adapters/kafka.py @@ -65,6 +65,7 @@ def __init__( rd_kafka_conf_options=None, debug: bool = False, poll_timeout: timedelta = timedelta(seconds=1), + broker_down_tolerance: timedelta = timedelta(milliseconds=500), rd_kafka_consumer_conf_options=None, rd_kafka_producer_conf_options=None, ): @@ -78,6 +79,8 @@ def __init__( set in this case since adapter will always replay from the last consumed offset. :param group_id_prefix - ( optional ) when not passing an explicit group_id, a prefix can be supplied that will be use to prefix the UUID generated for the group_id + :param broker_down_tolerance - how long all brokers may stay unreachable before the engine is shut down. librdkafka + reconnects on its own, so a single report is not evidence the brokers are gone. """ if group_id is not None and start_offset is not None: raise ValueError("start_offset is not supported when consuming with group_id") @@ -112,6 +115,7 @@ def __init__( "start_offset": start_offset.value if isinstance(start_offset, KafkaStartOffset) else start_offset, "max_threads": max_threads, "poll_timeout": poll_timeout, + "broker_down_tolerance": broker_down_tolerance, "rd_kafka_conf_properties": conf_properties, "rd_kafka_consumer_conf_properties": consumer_properties, "rd_kafka_producer_conf_properties": producer_properties, diff --git a/csp/tests/adapters/conftest.py b/csp/tests/adapters/conftest.py index 280cb3603..229fd4626 100644 --- a/csp/tests/adapters/conftest.py +++ b/csp/tests/adapters/conftest.py @@ -7,8 +7,9 @@ @pytest.fixture(scope="module", autouse=True) def kafkabroker(): - # Defined in ci/kafka/docker-compose.yml - return "localhost:9092" + # Defined in ci/kafka/docker-compose.yml. Not "localhost": that also resolves to ::1, where + # the container publishes nothing, and librdkafka alternates between the resolved addresses. + return "127.0.0.1:9092" @pytest.fixture(scope="module", autouse=True) diff --git a/csp/tests/adapters/test_kafka.py b/csp/tests/adapters/test_kafka.py index dd280bf92..0724d2162 100644 --- a/csp/tests/adapters/test_kafka.py +++ b/csp/tests/adapters/test_kafka.py @@ -390,7 +390,12 @@ def graph_sub(): with pytest.raises(RuntimeError): csp.run(graph_sub, starttime=utc_now(), endtime=timedelta(seconds=2), realtime=True) - kafkaadapter2 = KafkaAdapterManager(**kafkaadapterkwargs) + kafkaadapter2 = KafkaAdapterManager( + **kafkaadapterkwargs, + # The producer holds a message for a topic it cannot find until this elapses, on the + # assumption the topic is about to be created; the 30s default dominates the test + rd_kafka_producer_conf_options={"topic.metadata.propagation.max.ms": "2000"}, + ) def graph_pub(): msg_mapper = RawTextMessageMapper() @@ -398,7 +403,7 @@ def graph_pub(): # With bug this would deadlock with pytest.raises(RuntimeError): - csp.run(graph_pub, starttime=utc_now(), endtime=timedelta(seconds=2), realtime=True) + csp.run(graph_pub, starttime=utc_now(), endtime=timedelta(seconds=5), realtime=True) @pytest.mark.skipif(not os.environ.get("CSP_TEST_KAFKA"), reason="Skipping kafka adapter tests") def test_invalid_broker(self, kafkaadapterkwargs):