Skip to content

kafka_franz / redpanda output: no way to produce a tombstone (null value) after a mapping #4851

Description

@LostHeir

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.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions