Summary
The redpanda and kafka_franz outputs (and the others built on FranzWriter) cannot write a Kafka tombstone, meaning a record with a null value, once a message has gone through a Bloblang mapping. The advice given in #1999 and #3180 (root = "") does not produce a tombstone. It produces a record with a zero-length value, and log compaction keeps that record.
That makes Connect unusable as a transactional-outbox / CDC relay onto compacted topics whenever the relay has to publish deletes.
Reproduction (Connect 4.110.0, local Redpanda)
Pipeline: postgres_cdc (outbox table with a nullable BYTEA value) → mapping → redpanda output. Records read back with confluent_kafka (msg.value()):
| Mapping when the source value is null |
Bytes on the topic |
msg.value() |
root = this.value (i.e. root = null) |
4 B, the text null |
b'null' |
root = "" / "".bytes() / content().slice(0,0) / "".decode("base64") |
0 B |
b'' |
for reference, rpk topic produce -Z |
null (length -1) |
None |
No mapping we could find produces the third row.
Why
- Bloblang
root = null calls SetStructured(nil), which serializes to []byte("null").
root = "" produces a non-nil, zero-length byte slice.
FranzWriter.messageBatchToFranzRecords does r.Value, err = msg.AsBytes() and has no way to set a nil value. The sarama kafka output does the same with sarama.ByteEncoder(msgBytes).
- franz-go (
kbin.AppendVarintBytes) writes length -1 only for a nil slice.
kafka_tombstone_message exists only as input metadata. Tombstones survive a pure passthrough (the input creates the message from a nil value), but any processor that sets the content loses them.
Impact
- Consumers that decode with Schema Registry fail on the 4-byte
null and dead-letter the record, since it is not Confluent-framed.
- With
root = "", compaction never removes the key. The empty record stays forever on cleanup.policy=compact topics, which matters for GDPR-style deletes.
- Consumers can't tell "deleted" apart from "empty payload" without an out-of-band convention.
Proposal
Add an optional, advanced, interpolated boolean field to the shared FranzWriter config (FranzWriterConfigFields), for example tombstone. When it resolves to true, the record is written with a nil value. Key, headers, timestamp and partition still apply. The default is unset, so existing behaviour doesn't change.
output:
redpanda:
topic: ${! @kafka_topic }
key: ${! @kafka_key }
tombstone: ${! @is_delete } # or ${! @kafka_tombstone_message } for kafka -> kafka
Implementation sketch in messageBatchToFranzRecords, right after the value is set:
if w.Tombstone != nil {
s, err := w.Tombstone.TryString(msg)
if err != nil {
return nil, fmt.Errorf("tombstone interpolation: %w", err)
}
isTombstone, err := strconv.ParseBool(s)
if err != nil {
return nil, fmt.Errorf("parse tombstone: %w", err)
}
if isTombstone {
r.Value = nil
}
}
Alternatives considered:
- A Bloblang-level way to set the content to nil. This is more general but a larger change in benthos core.
- Treating zero-length content as a tombstone by default. This would be a breaking change, because empty values are legitimate in Kafka.
Summary
The
redpandaandkafka_franzoutputs (and the others built onFranzWriter) cannot write a Kafka tombstone, meaning a record with a null value, once a message has gone through a Bloblang mapping. The advice given in #1999 and #3180 (root = "") does not produce a tombstone. It produces a record with a zero-length value, and log compaction keeps that record.That makes Connect unusable as a transactional-outbox / CDC relay onto compacted topics whenever the relay has to publish deletes.
Reproduction (Connect 4.110.0, local Redpanda)
Pipeline:
postgres_cdc(outbox table with a nullableBYTEA value) →mapping→redpandaoutput. Records read back withconfluent_kafka(msg.value()):msg.value()root = this.value(i.e.root = null)nullb'null'root = ""/"".bytes()/content().slice(0,0)/"".decode("base64")b''rpk topic produce -ZNoneNo mapping we could find produces the third row.
Why
root = nullcallsSetStructured(nil), which serializes to[]byte("null").root = ""produces a non-nil, zero-length byte slice.FranzWriter.messageBatchToFranzRecordsdoesr.Value, err = msg.AsBytes()and has no way to set a nil value. The saramakafkaoutput does the same withsarama.ByteEncoder(msgBytes).kbin.AppendVarintBytes) writes length-1only for anilslice.kafka_tombstone_messageexists only as input metadata. Tombstones survive a pure passthrough (the input creates the message from a nil value), but any processor that sets the content loses them.Impact
nulland dead-letter the record, since it is not Confluent-framed.root = "", compaction never removes the key. The empty record stays forever oncleanup.policy=compacttopics, which matters for GDPR-style deletes.Proposal
Add an optional, advanced, interpolated boolean field to the shared
FranzWriterconfig (FranzWriterConfigFields), for exampletombstone. When it resolves totrue, the record is written with a nil value. Key, headers, timestamp and partition still apply. The default is unset, so existing behaviour doesn't change.Implementation sketch in
messageBatchToFranzRecords, right after the value is set:Alternatives considered: