Answer to the Core Question
To keep the same Producer ID across Cloud Function invocations you must keep a single Kafka producer instance alive inside the same container. In the OpenWhisk runtime this is achieved by creating the producer outside the handler function, in a module‑level variable. The container will reuse that object for every warm invocation, so the broker sees the same PID and idempotence suppresses duplicates. If the container is recycled or a new instance is spun up, a new PID is generated and the duplicate suppression window closes.
Impact of max.in.flight.requests.per.connection
Enabling enable.idempotence forces the client to limit max.in.flight.requests.per.connection to 5. For most serverless workloads that process one message per invocation, the effect on latency is negligible. In high‑volume functions that send many messages in a single batch, the limit can become a throughput bottleneck. A common mitigation is to increase batch.size and linger.ms so that fewer, larger batches are sent, keeping the in‑flight count low while still achieving the desired throughput.
Interaction with Cloud Functions Retries
Cloud Functions will retry an entire invocation when it throws an unhandled exception. If the retry occurs on a different container, the producer is re‑created and the broker treats the message as new, so the idempotence guarantee is lost. To bridge this gap you should use deterministic message keys (e.g., a business transaction ID) and design the downstream consumer to be idempotent (UPSERT or a unique constraint). That way, even if the broker receives the same key twice, the downstream system will apply the update only once.
Practical Steps for Your Use‑Case
- Instantiate the Kafka producer at module scope and keep it in a global variable.
- Set
enable.idempotence=true and verify the broker logs show the same Producer ID for consecutive messages from the same container.
- Use a business‑level key for every message and document that downstream consumers must deduplicate on that key.
- If you need higher throughput, tune
batch.size, linger.ms, and optionally max.in.flight.requests.per.connection (within the 5‑request limit).
- Monitor Cloud Functions logs for retry events and correlate them with Event Streams offsets to confirm no duplicate writes.
Diagnostic Question
Do you currently call producer.send(...).get() synchronously in the function, or do you rely on the async callback? The choice influences whether the function may terminate before the broker has confirmed the send, which can affect retry behavior.