[
https://issues.apache.org/jira/browse/KAFKA-20797?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Gaurav Narula reassigned KAFKA-20797:
-------------------------------------
Assignee: Pritam Kumar
> KIP-1365: Transform Observability and Skipped Record Handling for Kafka
> Connect
> -------------------------------------------------------------------------------
>
> Key: KAFKA-20797
> URL: https://issues.apache.org/jira/browse/KAFKA-20797
> Project: Kafka
> Issue Type: Improvement
> Components: connect, kip
> Reporter: Pritam Kumar
> Assignee: Pritam Kumar
> Priority: Major
>
> Kafka Connect's Single Message Transform (SMT) framework allows users to
> build chains of transformations that process records between the connector
> and Kafka. Transforms can modify records, route them to different topics, or
> *drop them entirely* by returning {{null}} (a filter operation). When
> combined with the error tolerance framework ({{{}errors.tolerance=all{}}}),
> records can also be diverted to a Dead Letter Queue (DLQ) before reaching the
> connector.
> Today, the transform layer has two significant gaps: *no operational
> observability* and {*}no mechanism for sink connectors to learn about records
> they never received{*}. This KIP addresses both with two complementary
> enhancements.
> h3. Problem 1: No Per-Transform Operational Metrics
> Kafka Connect provides *no operational metrics* at the transform layer. The
> existing {{connector-transform-metrics}} group only exposes static metadata:
>
> ||Existing Metric||Type||Description||
> |{{transform-class}}|Gauge (string)|The class name of the transformation|
> |{{transform-version}}|Gauge (string)|The version of the transformation|
> This creates several operational pain points:
> # *No visibility into filtering* — When using the {{Filter}} transform (or
> any transform that drops records by returning {{{}null{}}}), operators have
> no way to know how many records are being filtered. The only workaround is
> subtracting task-level {{source-record-poll-total}} from
> {{source-record-write-total}} (or {{sink-record-read-total}} from
> {{{}sink-record-send-total{}}}), which provides only an aggregate across the
> entire transform chain.
> # *No per-transform throughput* — In a chain of 3-5 transforms, operators
> cannot determine the throughput contribution or bottleneck of each individual
> transform. If a transform has a bug that silently drops records, it goes
> undetected.
> # *No error attribution* — While {{task-error-metrics}} tracks aggregate
> errors, operators cannot determine which transform in a chain is causing
> failures without correlating with debug-level logs.
> # *Capacity planning gap* — Without per-transform record counts, operators
> cannot make informed decisions about whether a transform is cost-effective
> (e.g., "is this filter dropping 90% of records, or 1%?").
> The consumer group protocol and other Kafka subsystems have rich
> per-component metrics. The transform layer is a notable gap.
> h3. Problem 2: Offset Gap for Sink Connectors When Records Are Skipped
> The {{SinkTask.preCommit()}} API allows sink connectors to control which
> offsets are committed to Kafka. Many production connectors override this
> method to implement exactly-once or at-least-once delivery guarantees with
> external systems. These connectors track offsets based on records they
> receive in {{put()}} and only return offsets from {{preCommit()}} after
> durably writing to the external system.
> When transforms or error tolerance drop records before {{{}put(){}}}, these
> connectors have *no visibility* into the dropped records. The result:
> # *Offset gap* — The consumer position advances past the dropped records,
> but the sink's committed offset does not include them.
> # *Reprocessing after rebalance* — On rebalance, the consumer seeks to the
> last committed offset and re-fetches the dropped records. The transforms drop
> them again, and the cycle repeats.
> # *Infinite reprocessing loop* — If all records for a partition are dropped,
> the sink returns an empty map from {{{}preCommit(){}}}, the framework skips
> commit entirely, and no progress is ever made for that partition.
> This problem affects widely-used connectors in production:
>
> ||Connector||Offset Strategy||Impact||
> |*S3 Sink*|Returns offsets only after flushing to S3 (rotation interval /
> record count)|Filtered records create offset gaps. Rebalance reprocesses all
> records since last S3 flush|
> |IcebergSinkConnector|Workers commit offset only after coordinator
> signal|Trailing filtered records never committed. Gap grows with each poll
> cycle|
> |*Any connector* using {{preCommit()}} returning {{{}}}|Full offset
> opt-out|If all records filtered, partition offset permanently frozen|
> h3. How the Two Problems Relate
> The two problems form a *diagnose → fix* pair:
> |Per-Transform Metrics (Part A) Skipped Record Notification (Part B)
> ───────────────────────────── ──────────────────────────────────────
> DIAGNOSE: "Which transform is FIX: "The sink can now track offsets dropping
> records, and how many?" for records it never received." Operator sees
> transform-record- Connector overrides onRecordSkipped() filtered-total
> climbing on to incorporate skipped offsets into filterNull transform at
> 632/sec its preCommit() return value │ │ └─── Together, they close the loop
> ──────┘|
>
>
> Without Part A, operators cannot diagnose *where* records are being dropped.
> Without Part B, connectors cannot *act* on records they never received.
> Together, they provide complete observability and correctness for the
> transform pipeline.
> The use of SMTs and Predicates in production has grown significantly since
> their introduction. Common patterns that trigger these problems:
> * *{{Filter}}* *SMT* — Drop records not matching criteria (e.g., filter by
> record type, tenant, region)
> * *{{RegexRouter}}* *+* *{{Predicate}}* — Route and filter based on content
> * *{{errors.tolerance=all}}* — Tolerate conversion errors, send to DLQ
> * *{{HeaderFilter}}* *+* *{{Predicate}}* — Drop records missing required
> headers
> * *Multi-tenant pipelines* — Filter per-tenant records across shared topics
> h4. *Success Criteria*
> ||Area||Criterion||
> |*Observability*|Operators can query per-transform record-in, record-out, and
> record-filtered counts via JMX|
> |*Observability*|Operators can compute per-transform filter rates for
> alerting and dashboards|
> |*Correctness*|Sink connectors that override {{preCommit()}} can track
> offsets for all consumed records, including those dropped before {{put()}}|
> |*Safety*|Zero risk of data loss — the framework never overrides the sink's
> offset decisions|
> |*Compatibility*|Full backward compatibility — existing connectors work
> identically without code changes|
> |*Performance*|Zero performance regression for connectors with no transforms
> configured|
> |*Universality*|Metrics are available for all transforms (built-in and
> custom) without any transform code changes|
> |*Simplicity*|Connectors need minimal code changes to incorporate skipped
> record awareness|
--
This message was sent by Atlassian Jira
(v8.20.10#820010)