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",
          access_key_id="test", secret_access_key="test")

There is no implicit credential fallback: either export the standard AWS_ACCESS_KEY_ID/AWS_SECRET_ACCESS_KEY env vars (or a profile) or pass access_key_id/secret_access_key/session_token explicitly. 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.1.tar.gz (767.6 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.1-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.1-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.1-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.1-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.1-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.1.tar.gz.

File metadata

  • Download URL: kcl_rs-0.1.1.tar.gz
  • Upload date:
  • Size: 767.6 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.1.tar.gz
Algorithm Hash digest
SHA256 bff6217a8c5cacbc19b73c3fdd1829995873ce633a519c682adb880dc1a50a93
MD5 dbe4f49b703729665b6a5766b162b326
BLAKE2b-256 dace898d2c733d88acffa2c8ea20f847aec091fe6370977fae176778d9e53718

See more details on using hashes here.

Provenance

The following attestation bundles were made for kcl_rs-0.1.1.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.1-cp39-abi3-musllinux_1_2_x86_64.whl.

File metadata

File hashes

Hashes for kcl_rs-0.1.1-cp39-abi3-musllinux_1_2_x86_64.whl
Algorithm Hash digest
SHA256 07d7a3115bec010fd1baa7808a5ea11a2a77d751f7fab45e86c98ea94fb084e4
MD5 ac3c1c54dabff3867364a72df7aa2c51
BLAKE2b-256 d6ad65727a554dfbd54b33d187de422931c0147e58b69971d2f3c99476098cd5

See more details on using hashes here.

Provenance

The following attestation bundles were made for kcl_rs-0.1.1-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.1-cp39-abi3-musllinux_1_2_aarch64.whl.

File metadata

File hashes

Hashes for kcl_rs-0.1.1-cp39-abi3-musllinux_1_2_aarch64.whl
Algorithm Hash digest
SHA256 60a1368134003c93e4c0ef680bc2e6c3ddb983a84e74dbaf3c8531876bf879ef
MD5 f5f2c3e08dbf83f05d8aeb85ceff6d05
BLAKE2b-256 43014f83129f0211515647d8e7e0bf4ef2107a6020097453eed4375a8e5d1aed

See more details on using hashes here.

Provenance

The following attestation bundles were made for kcl_rs-0.1.1-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.1-cp39-abi3-manylinux_2_28_x86_64.whl.

File metadata

File hashes

Hashes for kcl_rs-0.1.1-cp39-abi3-manylinux_2_28_x86_64.whl
Algorithm Hash digest
SHA256 7855fdfd54952c205f1cec783c74ab3ce02e9113ec96be3464f73939930062f5
MD5 63eceab558a6a4009be4edf737dff316
BLAKE2b-256 e1b3cf80a81261f8a4000201006cc1daa257c22fc5835a0447c2b0b348622c74

See more details on using hashes here.

Provenance

The following attestation bundles were made for kcl_rs-0.1.1-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.1-cp39-abi3-manylinux_2_28_aarch64.whl.

File metadata

File hashes

Hashes for kcl_rs-0.1.1-cp39-abi3-manylinux_2_28_aarch64.whl
Algorithm Hash digest
SHA256 3fcb2d4f96fd920a40cee0d207721bad0fb4bbb990dc893b4af849d2e1169eb1
MD5 0a4e224df7961038e98ada09dcfad42a
BLAKE2b-256 7a68158b31bd18444dade7eca67f4a83677e14ae89b331648201c8d7ebb73a5c

See more details on using hashes here.

Provenance

The following attestation bundles were made for kcl_rs-0.1.1-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.1-cp39-abi3-macosx_11_0_arm64.whl.

File metadata

File hashes

Hashes for kcl_rs-0.1.1-cp39-abi3-macosx_11_0_arm64.whl
Algorithm Hash digest
SHA256 1f8b8ceabe9b39ab8223f4138ec05c10d138c6f46cc7911f72f6c5cad6f83ea9
MD5 e3aa2ca85c288fe5c384438c04b47e3c
BLAKE2b-256 41ba9820aba78be25b04529903391228915a02f472d017bf1729e6940b14e069

See more details on using hashes here.

Provenance

The following attestation bundles were made for kcl_rs-0.1.1-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