Skip to main content

Rust-powered Python bindings for the Amazon Kinesis Client Library

Project description

kcl-rs

A Rust port of the Amazon Kinesis Client Library (KCL) with native Python bindings — the KCL v3 consumer framework without a JVM and without the MultiLangDaemon.

Like the Java KCL, this is the framework that turns "read from a Kinesis stream" into a production consumer fleet: it coordinates shard leases between workers through a DynamoDB table, checkpoints progress, follows resharding (parents before children), balances load across workers (KCL v3 variance-based assignment, including CPU-informed rebalancing), retrieves records via polling or enhanced fan-out, and publishes CloudWatch (or OpenTelemetry) metrics. The port is behaviorally faithful: same lease algorithms, same DynamoDB table schema, same lifecycle semantics.

Two crates in one workspace:

Crate What it is
kcl/ The core library — pure Rust, tokio + aws-sdk-*.
kcl-python/ PyO3 bindings, published to PyPI as kcl-rs (import name kcl_rs; the API mirrors amazon_kclpy).

Python

Install

pip install kcl-rs

Wheels are abi3 (one wheel per platform, CPython ≥ 3.9) for Linux (manylinux + musllinux, x86_64 + aarch64) and macOS (Apple Silicon). The package is fully typed (PEP 561): mypy --strict works out of the box, and your editor gets docstrings and signatures for everything.

Quick start

Subclass RecordProcessorBase, hand the class to a Scheduler, and run it. Each processor instance owns one shard.

from kcl_rs import RecordProcessorBase, Scheduler, CheckpointError


class MyProcessor(RecordProcessorBase):
    def initialize(self, initialize_input):
        print(f"starting on shard {initialize_input.shard_id}")

    def process_records(self, process_records_input):
        for record in process_records_input.records:
            handle(record.binary_data)          # raw bytes, already de-aggregated
        try:
            process_records_input.checkpointer.checkpoint()
        except CheckpointError as e:
            if e.value == "ThrottlingException":
                ...                             # back off and retry

    def lease_lost(self, lease_lost_input):
        pass                                    # lease taken by another worker

    def shard_ended(self, shard_ended_input):
        shard_ended_input.checkpointer.checkpoint()   # required to complete the shard

    def shutdown_requested(self, shutdown_requested_input):
        shutdown_requested_input.checkpointer.checkpoint()


scheduler = Scheduler(
    stream_name="my-stream",
    application_name="my-app",                  # also the lease-table name
    record_processor_factory=MyProcessor,       # called once per shard
    region="us-east-1",
)
scheduler.run()   # blocks; Ctrl-C triggers a graceful shutdown

scheduler.start() is the non-blocking variant (the KCL runs on background threads); call scheduler.shutdown() — from any thread — to stop either mode gracefully.

Where to start reading the stream

Scheduler(..., initial_position="TRIM_HORIZON")                  # oldest available
Scheduler(..., initial_position="LATEST")                        # default: only new records
Scheduler(..., initial_position="AT_TIMESTAMP",
          timestamp=datetime(2024, 1, 1, tzinfo=timezone.utc))   # or epoch seconds

The initial position only applies to shards without a checkpoint; a restarted application always resumes from its lease table.

Multi-stream mode

One worker can consume several streams — pass serialized stream identifiers ("accountId:streamName:creationEpoch") instead of a stream name:

Scheduler(
    stream_identifiers=[
        "123456789012:orders:1700000000",
        "123456789012:events:1700000001",
    ],
    application_name="my-app",
    record_processor_factory=MyProcessor,
)

Two-phase (prepared) checkpoints

For exactly-once-style side effects across failover, durably record a pending checkpoint first, commit it after the side effect:

prepared = checkpointer.prepare_checkpoint()
write_side_effect()
prepared.checkpoint()

LocalStack

Scheduler(..., endpoint_url="http://localhost:4566")

When endpoint_url is set, deterministic test credentials are supplied automatically. A runnable end-to-end sample lives at kcl-python/samples/.

Runtime model

The Rust KCL runs on its own tokio runtime on background threads with the GIL released — leases, retrieval, and metrics never contend with your Python code. The GIL is acquired only for the duration of each callback. This is the main practical difference from amazon_kclpy: no Java daemon, no stdin/stdout protocol, one process.


Rust

The core crate is a regular Rust library (not on crates.io yet — use a git or path dependency):

[dependencies]
kcl = { git = "https://github.com/localstack/kcl-rs", package = "kcl" }
tokio = { version = "1", features = ["rt-multi-thread", "macros", "signal"] }
aws-config = "1"
aws-sdk-kinesis = "1"
aws-sdk-dynamodb = "1"
aws-sdk-cloudwatch = "1"

Implement the (synchronous) ShardRecordProcessor trait and a factory, build the seven sub-configs with ConfigsBuilder, and drive the Scheduler. The full runnable version of this example is kcl/examples/single_stream.rs (cargo run --example single_stream):

use std::sync::Arc;

use kcl::common::ConfigsBuilder;
use kcl::coordinator::Scheduler;
use kcl::lifecycle::events::{
    InitializationInput, LeaseLostInput, ProcessRecordsInput, ShardEndedInput,
    ShutdownRequestedInput,
};
use kcl::processor::{
    ShardRecordProcessor, ShardRecordProcessorFactory, SingleStreamTracker, StreamTracker,
};

struct MyProcessor;

impl ShardRecordProcessor for MyProcessor {
    fn initialize(&mut self, input: InitializationInput) {
        println!("initializing on shard {:?}", input.shard_id());
    }

    fn process_records(&mut self, input: ProcessRecordsInput) {
        for record in input.records().unwrap_or(&[]) {
            println!("record pk={:?} seq={:?}", record.partition_key(), record.sequence_number());
        }
        if let Some(checkpointer) = input.checkpointer() {
            checkpointer.checkpoint().expect("checkpoint failed");
        }
    }

    fn lease_lost(&mut self, _input: LeaseLostInput) {}

    fn shard_ended(&mut self, input: ShardEndedInput) {
        // Required: completing the shard lets its children be processed.
        input.checkpointer().checkpoint().expect("checkpoint at shard end failed");
    }

    fn shutdown_requested(&mut self, input: ShutdownRequestedInput) {
        let _ = input.checkpointer().checkpoint();
    }
}

struct MyFactory;

impl ShardRecordProcessorFactory for MyFactory {
    fn shard_record_processor(&self) -> Box<dyn ShardRecordProcessor + Send + Sync> {
        Box::new(MyProcessor)
    }
}

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    // Region/credentials/endpoint come from the standard AWS environment
    // (AWS_PROFILE, AWS_REGION, AWS_ENDPOINT_URL, ...).
    let aws = aws_config::load_defaults(aws_config::BehaviorVersion::latest()).await;

    let tracker: Arc<dyn StreamTracker + Send + Sync> =
        Arc::new(SingleStreamTracker::from_stream_name("my-stream"));
    let configs = ConfigsBuilder::new(
        tracker,
        "my-app", // application name = lease table + CloudWatch namespace
        aws_sdk_kinesis::Client::new(&aws),
        aws_sdk_dynamodb::Client::new(&aws),
        aws_sdk_cloudwatch::Client::new(&aws),
        "worker-1",
        Arc::new(MyFactory),
    );

    let scheduler = Scheduler::new(
        configs.checkpoint_config(),
        configs.coordinator_config(),
        configs.lease_management_config(),
        configs.lifecycle_config(),
        configs.metrics_config(),
        configs.processor_config(),
        configs.retrieval_config(),
    )?;

    // Graceful shutdown on Ctrl-C.
    let scheduler_for_signal = Arc::clone(&scheduler);
    tokio::spawn(async move {
        tokio::signal::ctrl_c().await.expect("ctrl-c handler");
        scheduler_for_signal.shutdown().await;
    });

    scheduler.run().await; // blocks until shutdown completes
    Ok(())
}

Things worth knowing:

  • The processor callbacks are synchronous (matching Java's void signatures); all AWS I/O underneath is async. Blocking briefly in a callback is fine — callbacks run on blocking-capable threads, never on a tokio worker.
  • Scheduler::new must be called inside a tokio runtime (collaborators capture the runtime handle) but performs no network I/O; the first I/O happens when the scheduler runs.
  • Each of the seven config structs has fluent setters for the knobs you'd tune in Java (failover_time_millis, max_records, idle times, retrieval mode — polling vs enhanced fan-out — and so on). See the common, leases, and retrieval module docs.
  • Everything scales the same way as Java KCL: run more worker processes with the same application_name and leases redistribute automatically.

Consuming from LocalStack

Both languages honor the standard AWS environment, so no code changes:

AWS_ENDPOINT_URL=http://localhost:4566 AWS_REGION=us-east-1 \
AWS_ACCESS_KEY_ID=test AWS_SECRET_ACCESS_KEY=test \
  cargo run --example single_stream

Development

cargo test -p kcl                        # unit suite (~1,300 tests, hermetic)
cargo test -p kcl -- --include-ignored   # + integration tests (needs AWS or LocalStack)
cargo clippy --workspace --all-targets   # expect zero warnings

# Python bindings (tests embed a Python interpreter):
LD_LIBRARY_PATH=$(python3 -c 'import sysconfig; print(sysconfig.get_config_var("LIBDIR"))') \
  cargo test -p kcl-python

# Build the wheel:
cd kcl-python && maturin build --release

No system prerequisites beyond Rust and Python — the KPL protobuf wire format is decoded by hand-written code (no protoc), and release builds are size-optimized (stripped, thin LTO).

Porting conventions, subsystem status, and the decisions log live in PORTING.md; deliberately-skipped Java tests are documented in TEST-PARITY.md.

License

Apache-2.0, same as the upstream Amazon Kinesis Client Library.

Project details


Download files

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

Source Distribution

kcl_rs-0.1.0.tar.gz (760.7 kB view details)

Uploaded Source

Built Distributions

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

kcl_rs-0.1.0-cp39-abi3-musllinux_1_2_x86_64.whl (7.1 MB view details)

Uploaded CPython 3.9+musllinux: musl 1.2+ x86-64

kcl_rs-0.1.0-cp39-abi3-musllinux_1_2_aarch64.whl (6.7 MB view details)

Uploaded CPython 3.9+musllinux: musl 1.2+ ARM64

kcl_rs-0.1.0-cp39-abi3-manylinux_2_28_x86_64.whl (6.8 MB view details)

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

kcl_rs-0.1.0-cp39-abi3-manylinux_2_28_aarch64.whl (6.5 MB view details)

Uploaded CPython 3.9+manylinux: glibc 2.28+ ARM64

kcl_rs-0.1.0-cp39-abi3-macosx_11_0_arm64.whl (6.3 MB view details)

Uploaded CPython 3.9+macOS 11.0+ ARM64

File details

Details for the file kcl_rs-0.1.0.tar.gz.

File metadata

  • Download URL: kcl_rs-0.1.0.tar.gz
  • Upload date:
  • Size: 760.7 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? Yes
  • Uploaded via: twine/6.1.0 CPython/3.13.13

File hashes

Hashes for kcl_rs-0.1.0.tar.gz
Algorithm Hash digest
SHA256 4bf825fceedd87abd0dd41b3e0089fd0d8614d3b85a0e52f97ebe7a40ad8315f
MD5 195b0ae09032761e64061b1ca0d60ce3
BLAKE2b-256 18df68ab2f6534852f05457265be0fdd8ba0e1d87e3e08ff31911a2f02221147

See more details on using hashes here.

Provenance

The following attestation bundles were made for kcl_rs-0.1.0.tar.gz:

Publisher: ci.yml on localstack/kcl-rs

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

File details

Details for the file kcl_rs-0.1.0-cp39-abi3-musllinux_1_2_x86_64.whl.

File metadata

File hashes

Hashes for kcl_rs-0.1.0-cp39-abi3-musllinux_1_2_x86_64.whl
Algorithm Hash digest
SHA256 27d2a681789b11b08cbacadffa588a9f3976ab5a24f73afd878af2b35fb2b0ae
MD5 dc5519f007bdc38f96f0567cda7f5f88
BLAKE2b-256 d039c336658c4717280538d199bddd5b5f2aa6f5ec5f6f32752dd6114573edea

See more details on using hashes here.

Provenance

The following attestation bundles were made for kcl_rs-0.1.0-cp39-abi3-musllinux_1_2_x86_64.whl:

Publisher: ci.yml on localstack/kcl-rs

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

File details

Details for the file kcl_rs-0.1.0-cp39-abi3-musllinux_1_2_aarch64.whl.

File metadata

File hashes

Hashes for kcl_rs-0.1.0-cp39-abi3-musllinux_1_2_aarch64.whl
Algorithm Hash digest
SHA256 08f7198ea48317a3063439fd50c078a82b622aba4a400f0e4938483c4dd482c8
MD5 7ff8553501f503484b6f71eb3638efff
BLAKE2b-256 d0b38eafd3783e3536420ec762ba82a8953d67e2ef98aac2ca58a16b85d933f9

See more details on using hashes here.

Provenance

The following attestation bundles were made for kcl_rs-0.1.0-cp39-abi3-musllinux_1_2_aarch64.whl:

Publisher: ci.yml on localstack/kcl-rs

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

File details

Details for the file kcl_rs-0.1.0-cp39-abi3-manylinux_2_28_x86_64.whl.

File metadata

File hashes

Hashes for kcl_rs-0.1.0-cp39-abi3-manylinux_2_28_x86_64.whl
Algorithm Hash digest
SHA256 4ff76e4c3f9f16450726f28932d43e26d803eff452847ed2bd4e1c4127ffc1e8
MD5 e783e4bf42399f8e097c48322981de79
BLAKE2b-256 045e16f800c859778f3ad9d31bc6ab98414178ebf7529d73ba567fb015b4ed6a

See more details on using hashes here.

Provenance

The following attestation bundles were made for kcl_rs-0.1.0-cp39-abi3-manylinux_2_28_x86_64.whl:

Publisher: ci.yml on localstack/kcl-rs

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

File details

Details for the file kcl_rs-0.1.0-cp39-abi3-manylinux_2_28_aarch64.whl.

File metadata

File hashes

Hashes for kcl_rs-0.1.0-cp39-abi3-manylinux_2_28_aarch64.whl
Algorithm Hash digest
SHA256 2a59dcceecf8d0b1a185e8870dbd9ddf9efeac4c5bd0cf8e0273c50fbde57539
MD5 27311783b15a9a940b154f6c8e2f84c2
BLAKE2b-256 bfaed87207ec7ecc1515a63b5ca3c634c928194a1c05c6b804144b2229e2f87a

See more details on using hashes here.

Provenance

The following attestation bundles were made for kcl_rs-0.1.0-cp39-abi3-manylinux_2_28_aarch64.whl:

Publisher: ci.yml on localstack/kcl-rs

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

File details

Details for the file kcl_rs-0.1.0-cp39-abi3-macosx_11_0_arm64.whl.

File metadata

File hashes

Hashes for kcl_rs-0.1.0-cp39-abi3-macosx_11_0_arm64.whl
Algorithm Hash digest
SHA256 39a791f3932c3ca6f7cb5ced32387d6dac2d1867675bdeaf2c7f206c3c36e524
MD5 5c188511989cb4cfd9c797a19411973a
BLAKE2b-256 22696805f5e00b4158a9f7eae9dbcc964357eda044ad48544848ea39302b2bbc

See more details on using hashes here.

Provenance

The following attestation bundles were made for kcl_rs-0.1.0-cp39-abi3-macosx_11_0_arm64.whl:

Publisher: ci.yml on localstack/kcl-rs

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

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page