Skip to main content

hgraph-kafka

C++-first Kafka services for hgraph, implementing the contract in hgraph RFC 0015. The extension uses librdkafka's C API and exposes the same service model to native C++ and Python graphs.

One path-bound, multi-interface service_impl owns the Kafka clients for a configuration. Subscriptions, publish requests, explicit commits, and events all bind to that service instance. Graph output reaches Kafka through the service's sink inputs; records, delivery reports, and events re-enter the root graph through bounded push sources. Kafka clients and worker threads are created on graph start and stopped with the graph.

The public record and configuration shapes are hgraph compound scalars. Kafka headers preserve order, duplicates, null values, and empty byte strings.

Native C++

The installed package exports hgraph::kafka:

#include <hgraph/kafka/service.h>
#include <hgraph/kafka/value_builders.h>
#include <hgraph/lib/std/operators/conversion.h>
#include <hgraph/lib/std/operators/registration.h>

using namespace hgraph;
using namespace hgraph::kafka;

struct KafkaGraph {
    static constexpr auto name = "kafka_graph";

    static void compose(Wiring &w) {
        const auto path = service::path("primary");
        register_service(
            w, path,
            service_config().bootstrap_servers({Str{"localhost:9092"}}).build());

        auto key = wire<stdlib::const_, TS<KafkaSubscriptionKey>>(
            w, subscription_key()
                   .topics({Str{"orders"}})
                   .group_id(Str{"orders-worker"})
                   .build());
        auto subscription = subscribe(w, path, key);

        auto record = wire<stdlib::const_, TS<KafkaProduceRecord>>(
            w, make_produce_record(Bytes{"ready"}));
        auto delivery = publish(
            w, path, publish_request(w, Str{"status"}, record));

        auto cursor = wire<stdlib::getattr_, TS<KafkaCursor>>(
            w, subscription, Str{"cursor"});
        commit(w, path, cursor);
        auto event = events(w, path);
    }
};

KafkaSubscriptionOutput provides the record and its matching next-offset cursor on the same graph tick, plus subscription state. A cursor is accepted only while its subscription identity, assignment generation, and partition remain live. Commits are monotonic per assigned partition.

Python

The Python authoring surface lowers to the same native service:

import hgraph as hg
import hgraph_kafka as kafka

@hg.graph
def app():
    kafka.register_kafka_service(
        kafka.KafkaServiceConfig.from_bootstrap_servers(
            ["localhost:9092"], client_id="orders-worker"
        ),
        path="primary",
    )
    key = kafka.KafkaSubscriptionKey(
        topics=("orders",),
        group_id="orders-worker",
        start_position=kafka.KafkaStartPosition.committed(),
    )
    subscription = kafka.kafka_subscribe(
        hg.const(key, tp=hg.TS[kafka.KafkaSubscriptionKey]),
        path="primary",
    )
    kafka.kafka_commit(subscription["cursor"], path="primary")

The extension wheel also owns the released import path hgraph.adaptors.kafka. Existing message_publisher, message_subscriber, KafkaMessage, and register_kafka_adaptor imports therefore continue to work without making the core hgraph package depend on this extension.

Recovery and simulation

Subscriptions support explicit topic, pattern, or partition selection; group or independent assignment; earliest, latest, committed, timestamp, explicit, and graph-start positions; snapshot, timestamp, and explicit stop boundaries; key filters; deterministic timestamp/topic/partition/offset replay; and explicit or graph-delivery commits.

Simulation is intentionally limited to bounded, record-time recovery. The consumer preloads the finite replay and schedules records at deterministic graph times. Publish, commit, unbounded asynchronous input, and OnGraphDelivery commit mode are rejected in simulation rather than silently changing their semantics.

Build and test

This is a first-party extension in the hgraph monorepo. It remains a separate CMake package and Python distribution: the top-level core package does not link librdkafka or install these modules.

For an in-tree native development build from the repository root:

cmake -S . -B build-kafka \
  -DHGRAPH_BUILD_KAFKA_EXTENSION=ON \
  -DBUILD_TESTING=ON
cmake --build build-kafka --parallel
ctest --test-dir build-kafka --output-on-failure

The extension can still be configured independently against an installed hgraph SDK:

cmake -S . -B build -DCMAKE_PREFIX_PATH=/path/to/hgraph/install
cmake --build build --parallel
ctest --test-dir build --output-on-failure

Build its separately deployable ABI3 wheel from the repository root after making the matching hgraph SDK discoverable through CMAKE_PREFIX_PATH:

CMAKE_PREFIX_PATH=/path/to/hgraph/sdk \
  uv build --wheel --package hgraph-kafka --python 3.12

The deterministic suite uses librdkafka's mock cluster and the extension fake transport. To include a real broker round trip, provide a clean topic:

HGRAPH_KAFKA_INTEGRATION_BOOTSTRAP=localhost:9092 \
HGRAPH_KAFKA_INTEGRATION_TOPIC=hgraph-kafka-integration \
ctest --test-dir build --output-on-failure

Wheel builds require the SDK installed by a stable-ABI hgraph wheel. The extension rejects an SDK that links Python::Python, because that would pin the nominal ABI3 module to the build interpreter.

Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

hgraph_kafka-0.8.4.tar.gz (80.1 kB view details)

Uploaded Source

Built Distributions

If you're not sure about the file name format, learn more about wheel file names.

hgraph_kafka-0.8.4-cp312-abi3-win_amd64.whl (3.0 MB view details)

Uploaded CPython 3.12+Windows x86-64

hgraph_kafka-0.8.4-cp312-abi3-manylinux_2_28_x86_64.whl (1.7 MB view details)

Uploaded CPython 3.12+manylinux: glibc 2.28+ x86-64

hgraph_kafka-0.8.4-cp312-abi3-macosx_15_0_arm64.whl (3.8 MB view details)

Uploaded CPython 3.12+macOS 15.0+ ARM64

File details

Details for the file hgraph_kafka-0.8.4.tar.gz.

File metadata

  • Download URL: hgraph_kafka-0.8.4.tar.gz
  • Upload date:
  • Size: 80.1 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? Yes
  • Uploaded via: twine/7.0.0 CPython/3.13.14

File hashes

Hashes for hgraph_kafka-0.8.4.tar.gz
Algorithm Hash digest
SHA256 632aa98a6c234368547a569d095f3f131b03c07c836ccacc1bdc806bdc5ae36e
MD5 7e372bbcfd8ae76f629cc9f2c7534b6d
BLAKE2b-256 1d36cdb0601b44ac43474a67d3e8c4ad1dde97c6291b18b8e78fd94186a7b1ab

See more details on using hashes here.

Provenance

The following attestation bundles were made for hgraph_kafka-0.8.4.tar.gz:

Publisher: release-wheels.yml on hhenson/hgraph

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

File details

Details for the file hgraph_kafka-0.8.4-cp312-abi3-win_amd64.whl.

File metadata

  • Download URL: hgraph_kafka-0.8.4-cp312-abi3-win_amd64.whl
  • Upload date:
  • Size: 3.0 MB
  • Tags: CPython 3.12+, Windows x86-64
  • Uploaded using Trusted Publishing? Yes
  • Uploaded via: twine/7.0.0 CPython/3.13.14

File hashes

Hashes for hgraph_kafka-0.8.4-cp312-abi3-win_amd64.whl
Algorithm Hash digest
SHA256 07edc8d0c9e356f9d640e64f7cece9a9137188480237a16b224cbdff8792639e
MD5 6dd1729aa4cfd61ee81caada9b4408c7
BLAKE2b-256 7e60ac82bb7dba943bcaa8cc519af9e47f55084c83453c72087a8a789687fd1c

See more details on using hashes here.

Provenance

The following attestation bundles were made for hgraph_kafka-0.8.4-cp312-abi3-win_amd64.whl:

Publisher: release-wheels.yml on hhenson/hgraph

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

File details

Details for the file hgraph_kafka-0.8.4-cp312-abi3-manylinux_2_28_x86_64.whl.

File metadata

File hashes

Hashes for hgraph_kafka-0.8.4-cp312-abi3-manylinux_2_28_x86_64.whl
Algorithm Hash digest
SHA256 b03325f9f467609d1d9b58d242f4c8b7be2f9d0f4d5d436c8feee430dd294016
MD5 1db22097cfc4024551785a3e57ad44a5
BLAKE2b-256 3fd9443e3f990366b8418f3341e7b6147633f814ec1f4db90e46d0925733f08c

See more details on using hashes here.

Provenance

The following attestation bundles were made for hgraph_kafka-0.8.4-cp312-abi3-manylinux_2_28_x86_64.whl:

Publisher: release-wheels.yml on hhenson/hgraph

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

File details

Details for the file hgraph_kafka-0.8.4-cp312-abi3-macosx_15_0_arm64.whl.

File metadata

File hashes

Hashes for hgraph_kafka-0.8.4-cp312-abi3-macosx_15_0_arm64.whl
Algorithm Hash digest
SHA256 7b62352c6894bbddb2df94b305e9ac89bc078e33be2d9b0d9f21d2cc88fa060c
MD5 7eb9c7e37efbf9ced5bb545f26ea2cd8
BLAKE2b-256 be0cdf8a57092bf8165a12a4b1504d74e9d9c25845a638adab7293a2cc697c34

See more details on using hashes here.

Provenance

The following attestation bundles were made for hgraph_kafka-0.8.4-cp312-abi3-macosx_15_0_arm64.whl:

Publisher: release-wheels.yml on hhenson/hgraph

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

Release history Release notifications | RSS feed

0.8.22

4 files

0.8.21

4 files

0.8.20

4 files

0.8.19

4 files

0.8.18

4 files

0.8.17

4 files

0.8.16

4 files

0.8.15

4 files

0.8.14

4 files

0.8.13

4 files

0.8.12

4 files

0.8.11

4 files

0.8.10

4 files

0.8.9

4 files

0.8.8

4 files

0.8.7

4 files

0.8.6

4 files

0.8.5

4 files

This release

0.8.4 This release

4 files

0.8.3

4 files

0.8.2

4 files

0.8.1

4 files

0.8.0

4 files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page