Skip to main content

amgi-kafka-event-source-mapping

amgi-kafka-event-source-mapping is an adaptor for AMGI applications to run in a Kafka event source mapped environment.

Installation

pip install amgi-kafka-event-source-mapping==0.47.0

Example

This example uses AsyncFast:

from dataclasses import dataclass

from amgi_kafka_event_source_mapping import KafkaEventSourceMappingHandler
from asyncfast import AsyncFast

app = AsyncFast()


@dataclass
class Order:
    item_ids: list[str]


@app.channel("orders")
async def orders(order: Order) -> None:
    # Makes an order
    ...


handler = KafkaEventSourceMappingHandler(app)

What it does

  • Converts Kafka batch events into AMGI message.receive events
  • Uses the Kafka topic name as the AMGI message address
  • Supports partial batch failures so only failed records are reported
  • Sends outbound messages to Kafka using an async producer
  • Outbound messages are sent via the same Kafka broker (bootstrap servers) that the records were received from
  • Optionally manages application startup and shutdown via AMGI lifespan

Record handling

  • Record values and keys are passed to your app as bytes
  • Kafka record headers become AMGI headers
  • Records are only acknowledged when your app emits message.ack
  • Records that emit message.nack or are not acknowledged are treated as failures

Nack handling

By default, records that are negatively acknowledged, or not acknowledged are logged:

handler = KafkaEventSourceMappingHandler(app, on_nack="log")

To fail the invocation when any record is nacked, configure the handler to raise an error instead:

handler = KafkaEventSourceMappingHandler(app, on_nack="error")

This is useful when running in environments where a failed invocation should trigger a retry, or alert.

When using this mode, handlers must be idempotent. Kafka event source mappings may re-deliver records after failures, restarts, or rebalances, and your application logic should be safe to execute more than once for the same record.

Lifespan

Lifespan support is enabled by default.

  • Startup runs once per Lambda execution environment
  • Shutdown is attempted when the environment is terminated

Shutdown handling relies on signal.SIGTERM, which is supported by Python 3.12 and later Lambda runtimes.

To use fully stateless, per-invocation behavior, disable lifespan:

handler = KafkaEventSourceMappingHandler(app, lifespan=False)

Contact

For questions or suggestions, please contact jack.burridge@mail.com.

License

Copyright 2026 AMGI

Download files

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

Source Distribution

amgi_kafka_event_source_mapping-0.47.0.tar.gz (5.5 kB view details)

Uploaded Source

Built Distribution

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

File details

Details for the file amgi_kafka_event_source_mapping-0.47.0.tar.gz.

File metadata

  • Download URL: amgi_kafka_event_source_mapping-0.47.0.tar.gz
  • Upload date:
  • Size: 5.5 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? Yes
  • Uploaded via: uv/0.12.5 {"installer":{"name":"uv","version":"0.12.5","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Ubuntu","version":"24.04","id":"noble","libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":true}

File hashes

Hashes for amgi_kafka_event_source_mapping-0.47.0.tar.gz
Algorithm Hash digest
SHA256 c532106c7fd43685127666d2a8410d87239cb5758ba2b5975f6bef1d4df0b153
MD5 b8bd50dc6d65eb01543e13d29b7ed0d9
BLAKE2b-256 fe253317d9250dc1ce9b528958c6537a33e804b1f6ba45cef97cd6b4438b4acd

See more details on using hashes here.

File details

Details for the file amgi_kafka_event_source_mapping-0.47.0-py3-none-any.whl.

File metadata

  • Download URL: amgi_kafka_event_source_mapping-0.47.0-py3-none-any.whl
  • Upload date:
  • Size: 6.8 kB
  • Tags: Python 3
  • Uploaded using Trusted Publishing? Yes
  • Uploaded via: uv/0.12.5 {"installer":{"name":"uv","version":"0.12.5","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Ubuntu","version":"24.04","id":"noble","libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":true}

File hashes

Hashes for amgi_kafka_event_source_mapping-0.47.0-py3-none-any.whl
Algorithm Hash digest
SHA256 65c49820cb0efaf22e13a002a7cb3a6a902e72cef08c1f46b292b84020578b20
MD5 4c80f25563f46b8fef6775dceb648674
BLAKE2b-256 c8abf4ac3d70901fa0cb74e2adbb44ecb3b159c2358e01a969718cdfff80dfcb

See more details on using hashes here.

Release history Release notifications | RSS feed

This release

0.47.0 This release

2 files

0.46.0

2 files

0.45.1

2 files

0.45.0

2 files

0.44.0

2 files

0.43.0

2 files

0.42.0

2 files

0.41.1

2 files

0.41.0

2 files

0.40.0

2 files

0.39.0

2 files

0.38.0

2 files

0.37.0

2 files

0.36.0

2 files

0.35.0

2 files

0.34.0

2 files

0.33.0

2 files

Supported by

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