How would you implement client-side (application-level) encryption for PII fields in Kafka messages, and what are the trade-offs?
answer
- Envelope: KMS DEK + wrapped DEK in headers
- AES-256-GCM, unique IV per message
- Serializer / interceptor wiring
- Key in cleartext, value ciphertext
- Breaks value-based compaction/filtering
basics
~20 sEncrypt the sensitive payload in the producer before send() and decrypt in the consumer after poll(), typically via a custom Serializer or interceptor using envelope encryption: a KMS issues a data key that encrypts the message, and the wrapped data key travels alongside. The broker stores only ciphertext.
solid answer
~50 sImplement encryption above Kafka so the broker never sees plaintext PII. The producer encrypts the message value (or specific fields) before sending; the consumer decrypts after polling. A clean pattern is envelope encryption: call a KMS (AWS KMS, Vault, GCP KMS) to get a data encryption key (DEK), encrypt the payload with the DEK using AES-GCM, and ship the KMS-wrapped DEK plus IV in the message (often in headers). This is wired via a custom Serializer/Deserializer or a ProducerInterceptor/ConsumerInterceptor so application code stays clean. Keep the partition key in cleartext so partitioning still works. Trade-offs: the broker can no longer inspect content, so log compaction by value semantics and any server-side filtering see only ciphertext; schema-registry validation must run before encryption; key rotation and crypto-shredding (deleting a key to render data unreadable) need design; and there is CPU/latency overhead. It does, however, defend PII even against broker operators.
code
java · 11 lines// ProducerInterceptor sketch: envelope-encrypt the value, keep key cleartext
public ProducerRecord<String, byte[]> onSend(ProducerRecord<String, byte[]> rec) {
DataKey dk = kms.generateDataKey("alias/kafka-pii"); // plaintext DEK + wrapped DEK
byte[] iv = randomNonce(12);
byte[] ciphertext = aesGcmEncrypt(dk.plaintext(), iv, rec.value(), rec.topic().getBytes());
Arrays.fill(dk.plaintext(), (byte) 0); // wipe DEK from memory
Headers h = rec.headers();
h.add("enc-dek", dk.wrapped());
h.add("enc-iv", iv);
return new ProducerRecord<>(rec.topic(), rec.partition(), rec.key(), ciphertext, h);
}go deeper
Understand the idea: encrypt in the producer, decrypt in the consumer, broker sees only ciphertext.
Explain envelope encryption (DEK + wrapped DEK), where to put the wrapped key (headers), and that the partition key stays cleartext.
Wire it via Serializer/interceptor, reason about AES-GCM/IV, KMS caching, and the impact on compaction, filtering, and schema validation.
Design org-wide key management, rotation, crypto-shredding for GDPR, field-level encryption with Schema Registry, and performance/throughput trade-offs.
## Why client-side encryption Kafka stores log segments as plaintext and disk encryption is decrypted by the running broker, so a broker operator or anyone reading through the broker can see message contents. For regulated **PII** (personally identifiable information — names, emails, SSNs), teams encrypt **in the client** so the broker only ever holds **ciphertext**. This is the strongest at-rest protection because trust never extends to the broker. ## Envelope encryption — the standard pattern Envelope encryption avoids sending raw data to a KMS for every message: 1. **Get a data key (DEK):** ask a **KMS** (Key Management Service — AWS KMS, HashiCorp Vault, GCP KMS, Azure Key Vault) for a symmetric **data encryption key**. The KMS returns the DEK in plaintext *and* a copy **wrapped (encrypted) by a key encryption key (KEK)** that never leaves the KMS. 2. **Encrypt the payload** locally with the DEK using an authenticated cipher like **AES-256-GCM** (GCM gives confidentiality + integrity via an auth tag, and needs a unique **IV**/nonce per message). 3. **Discard the plaintext DEK** from memory; **store/transmit the wrapped DEK + IV** alongside the ciphertext — typically in **Kafka record headers** (`record.headers()`), keeping the value field as pure ciphertext. 4. **Consumer side:** read the wrapped DEK from the header, call the KMS to **unwrap** it (the KMS uses the KEK), then AES-GCM-decrypt the payload. DEKs are usually **cached/reused** for many messages to avoid a KMS call per record, balancing security and throughput. ## Where to wire it - **Custom Serializer/Deserializer:** encryption happens inside `serialize()`/`deserialize()`. Clean but couples crypto to serialization. - **ProducerInterceptor / ConsumerInterceptor:** `onSend()` / `onConsume()` transform records; keeps serializers standard. - Field-level: encrypt only sensitive fields (e.g., the `email` field) so the rest stays queryable. Confluent's client-side field-level encryption with Schema Registry tags does exactly this. ## Keep these in cleartext - **The message key / partition key** — otherwise hashing changes per encryption and ordering/partitioning breaks. - **Routing headers** the broker or stream processors need. ## Trade-offs and edge cases - **No content inspection by the broker:** **log compaction** still works (it dedupes by *key*, which is cleartext) but you cannot compact or filter by *value content*; server-side filtering and ksqlDB transforms see only ciphertext. - **Schema validation order:** validate against Schema Registry **before** encrypting; the broker/registry can't validate ciphertext. - **Key rotation & crypto-shredding:** rotating the KEK re-wraps DEKs cheaply. For **GDPR erasure**, **crypto-shredding** — destroying the key for a subject so their ciphertext becomes permanently undecryptable — is a powerful complement to deletion (see retention/tombstones). - **Performance:** AES-GCM is hardware-accelerated (AES-NI) and cheap; the cost is mostly KMS round-trips, mitigated by DEK caching. - **Replay/integrity:** GCM's auth tag detects tampering; include context (e.g., topic) as AAD to bind ciphertext to its use. ## What NOT to do - Don't encrypt the partition key. Don't put the plaintext DEK in the message. Don't rely on the broker for any value-based logic once encrypted. Don't reuse an IV with the same key (GCM nonce reuse is catastrophic).
- Why use envelope encryption instead of calling the KMS to encrypt every message directly?Sending every payload to the KMS is slow, costly, and rate-limited. Envelope encryption fetches a data key (DEK) once, encrypts many messages locally with AES-GCM, and only stores the small KMS-wrapped DEK. The KMS only ever handles tiny keys, not bulk data.
- Why must the partition key stay in cleartext?Kafka chooses a partition by hashing the record key. If you encrypt the key, the hash changes (and is non-deterministic with random IVs), so messages for the same logical entity scatter across partitions, breaking ordering and key-based compaction.
saying these in an interview costs you the question
- Encrypting the partition/message key, breaking partitioning and ordering
- Storing the plaintext data key in the message
- Calling the KMS to encrypt the full payload for every record
- Claiming log compaction breaks entirely (it dedupes by cleartext key, not value)
- Reusing an IV/nonce with the same key under AES-GCM