Skip to main content

Kafka integration

Kafka can be the managed Command source, the delivery destination, or both. Nereus Delay preserves Kafka-specific resource and evidence semantics behind the common delayed-message lifecycle.

Kafka as a Command source​

One physical Kafka topic partition maps to one Delay Shard. Its record offset is the Source Position within that Shard Log. Source positions from different topics, route incarnations, or partitions are not globally comparable.

The Worker applies each accepted record to the Shard's RocksDB database before committing source progress. Recovery uses the durable applied position and replays without acknowledging records that have not been synchronously materialized.

The guarded source binds the configured Kafka cluster, topic, native topic ID, and partition at the actual Fetch boundary. Recreating a topic with the same name produces a different resource incarnation and must fail closed rather than feed a Shard from an unintended log.

Kafka as a destination​

For ordinary managed delivery, Kafka uses actionAt = deliverAt: Nereus Delay does not call the target producer before the trusted time boundary permits it. Backlog, retry, or target throttling can make delivery later.

The destination binding fixes the target profile and physical partition when the Schedule is applied. Reschedule changes time and generation, not the destination or payload. Ordering claims are limited to the registered single-source-partition and single-target-partition domain.

Delivery evidence​

The baseline capability is bounded at-least-once. A lost response after producer ownership can become UNCERTAIN, because the target append may have succeeded. Nereus Delay does not turn that window into a false failure.

A stronger Kafka transactional-receipt profile is opt-in. It is valid only when the configured target, receipt resource, transactional identity, credential binding, and resource-incarnation checks all remain certified. If those prerequisites drift, the affected Destination Lane becomes not ready; the system does not silently downgrade a declared guarantee.

Operational checklist​

  • Treat the native topic ID, not only the topic name, as part of source and target identity.
  • Keep source acknowledgement after the RocksDB synchronous apply boundary.
  • Monitor source lag separately from destination Lane lag.
  • Retain delayMessageId + generation for application-level duplicate handling.
  • Treat transactional receipt as a capability profile, not as a universal Kafka property.

See delivery guarantees and uncertainty and Destination Lane isolation.