# Stream Processing

**URL:** https://forum.confluent.io/c/stream-processing/43.md

[Latest](https://forum.confluent.io/latest.md) · [Categories](https://forum.confluent.io/categories.md)

---

## [About the Stream Processing category](https://forum.confluent.io/t/about-the-stream-processing-category/7944)

<div class="topic-metadata">

**Author:** [@Joe](https://forum.confluent.io/u/Joe)\
**Replies:** 0

</div>

All things Stream Processing!

---

## [Issue after upgrading to CMF 2.4.2/CFK 3.3.0: FlinkApplication resolution fails cross-namespace FlinkEnvironment](https://forum.confluent.io/t/issue-after-upgrading-to-cmf-2-4-2-cfk-3-3-0-flinkapplication-resolution-fails-cross-namespace-flinkenvironment/38516)

<div class="topic-metadata">

**Author:** [@TinyToons](https://forum.confluent.io/u/TinyToons)\
**Replies:** 5\
**Last updated:** [30 September 2026 09:52 UTC](https://forum.confluent.io/t/issue-after-upgrading-to-cmf-2-4-2-cfk-3-3-0-flinkapplication-resolution-fails-cross-namespace-flinkenvironment/38516 "2026-09-30T09:52:39Z")

</div>

Hi everyone, After upgrading Confluent Manager for Flink (CMF) from 2.3.2 to 2.4.2 and Confluent For Kubernetes (CFK) from 3.2.1 to 3.3.0, CFK fails to reconcile an existing FlinkApplication. It seems to look for the F…

---

## [CFK 3.3.0 - FlinkEnvironment is never reconciled, operator logs cmfDay2OpsEnabled=false](https://forum.confluent.io/t/cfk-3-3-0-flinkenvironment-is-never-reconciled-operator-logs-cmfday2opsenabled-false/38470)

<div class="topic-metadata">

**Author:** [@LHP02052003](https://forum.confluent.io/u/LHP02052003)\
**Replies:** 1\
**Last updated:** [28 September 2026 13:00 UTC](https://forum.confluent.io/t/cfk-3-3-0-flinkenvironment-is-never-reconciled-operator-logs-cmfday2opsenabled-false/38470 "2026-09-28T13:00:28Z")

</div>

Hi everyone, I’m evaluating Confluent Manager for Apache Flink (CMF) together with Confluent for Kubernetes (CFK) on OpenShift. I’m encountering an issue where a FlinkEnvironment resource is created successfully but is…

---

## [ksqlDB JSON\_SR with existing KEY\_SCHEMA\_ID fails on structured key: Incompatible schema of type JSON / id.compatibility.strict](https://forum.confluent.io/t/ksqldb-json-sr-with-existing-key-schema-id-fails-on-structured-key-incompatible-schema-of-type-json-id-compatibility-strict/38374)

<div class="topic-metadata">

**Author:** [@Ethan2029](https://forum.confluent.io/u/Ethan2029)\
**Replies:** 0\
**Last updated:** [17 May 2026 11:45 UTC](https://forum.confluent.io/t/ksqldb-json-sr-with-existing-key-schema-id-fails-on-structured-key-incompatible-schema-of-type-json-id-compatibility-strict/38374 "2026-05-17T11:45:44Z")

</div>

Hi Confluent community, I am looking for advice on a ksqlDB / JSON\_SR / Schema Registry issue. Context We have the following flow: Source topic: raw JSON, no Schema Registry ksqlDB: transforms the payload into the ta…

---

## [Kafka Streams EOS - Producer fenced](https://forum.confluent.io/t/kafka-streams-eos-producer-fenced/38325)

<div class="topic-metadata">

**Author:** [@roookeee](https://forum.confluent.io/u/roookeee)\
**Replies:** 4\
**Last updated:** [16 April 2026 09:08 UTC](https://forum.confluent.io/t/kafka-streams-eos-producer-fenced/38325 "2026-04-16T09:08:49Z")

</div>

I am currently trying to get rid of the following error in our EOS-configured Kafka Streams Spring Boot application: org.apache.kafka.common.errors.InvalidProducerEpochException: Producer attempted to produce with an o…

---

## [Kafka Streams FK-join: force result timestamp to always come from main (left) flow](https://forum.confluent.io/t/kafka-streams-fk-join-force-result-timestamp-to-always-come-from-main-left-flow/38286)

<div class="topic-metadata">

**Author:** [@foal](https://forum.confluent.io/u/foal)\
**Replies:** 5\
**Last updated:** [20 March 2026 17:23 UTC](https://forum.confluent.io/t/kafka-streams-fk-join-force-result-timestamp-to-always-come-from-main-left-flow/38286 "2026-03-20T17:23:56Z")

</div>

Hi, we use Kafka Streams 4.1.1 and need guidance on FK-join timestamp behaviour. Concrete example (generalised): Main flow record (authoritative entity): timestamp T\_main=1000 Enrichment flow record (auxiliary update)…

---

## [KTable-KTable Foreign-Key LEFT JOIN: Discrepancy between documentation and behavior when FK is NULL](https://forum.confluent.io/t/ktable-ktable-foreign-key-left-join-discrepancy-between-documentation-and-behavior-when-fk-is-null/38300)

<div class="topic-metadata">

**Author:** [@foal](https://forum.confluent.io/u/foal)\
**Replies:** 1\
**Last updated:** [20 March 2026 17:22 UTC](https://forum.confluent.io/t/ktable-ktable-foreign-key-left-join-discrepancy-between-documentation-and-behavior-when-fk-is-null/38300 "2026-03-20T17:22:19Z")

</div>

KTable-KTable Foreign-Key LEFT JOIN: Discrepancy between documentation and behavior when FK is NULL Summary I’m observing a discrepancy between the documented semantics table and the actual runtime behavior of KTable-KTa…

---

## [Unable to migrate metadata from zookeeper to KRaft](https://forum.confluent.io/t/unable-to-migrate-metadata-from-zookeeper-to-kraft/38287)

<div class="topic-metadata">

**Author:** [@Agaikwad](https://forum.confluent.io/u/Agaikwad)\
**Replies:** 0\
**Last updated:** [13 March 2026 13:07 UTC](https://forum.confluent.io/t/unable-to-migrate-metadata-from-zookeeper-to-kraft/38287 "2026-03-13T13:07:51Z")

</div>

Using Confluent kafka 7.9.4-1. Following the documentation of migration from zookeeper to KRaft i completed setting up my controller and its able to communicate with my broker . but the migration is not getting complet…

---

## [Kafka Streams and Schemaregistry interaction, multiple teams, languages](https://forum.confluent.io/t/kafka-streams-and-schemaregistry-interaction-multiple-teams-languages/38277)

<div class="topic-metadata">

**Author:** [@lukas\_dev](https://forum.confluent.io/u/lukas_dev)\
**Replies:** 2\
**Last updated:** [6 March 2026 13:16 UTC](https://forum.confluent.io/t/kafka-streams-and-schemaregistry-interaction-multiple-teams-languages/38277 "2026-03-06T13:16:45Z")

</div>

Hello! We use Confluent and Schemaregistry, with protos. There is an upstream team working in Dotnet, which makes schema evolution progress. I work in the downstream BI team, working in Java. We consume from their top…

---

## [Kafka Streams - bizarre data loss](https://forum.confluent.io/t/kafka-streams-bizarre-data-loss/38261)

<div class="topic-metadata">

**Author:** [@roookeee](https://forum.confluent.io/u/roookeee)\
**Replies:** 8\
**Last updated:** [5 March 2026 21:13 UTC](https://forum.confluent.io/t/kafka-streams-bizarre-data-loss/38261 "2026-03-05T21:13:52Z")

</div>

We are currently investigating a data loss in our Kafka Streams application. Both the state store and outbound topic message got lost for a single message (which we detected by pure chance, and yes, only one message got…

---

## [CompleteBatch Infinite Loop With BatchSize Greater Than 1](https://forum.confluent.io/t/completebatch-infinite-loop-with-batchsize-greater-than-1/37640)

<div class="topic-metadata">

**Author:** [@khamburg](https://forum.confluent.io/u/khamburg)\
**Replies:** 6\
**Last updated:** [2 February 2026 21:13 UTC](https://forum.confluent.io/t/completebatch-infinite-loop-with-batchsize-greater-than-1/37640 "2026-02-02T21:13:39Z")

</div>

I have a Kafka Streams app and I’m implementing a production exception handler for messages that are too large for the destination topic. I noticed that if return FAIL from the handler and then restart my app, it will go…

---

## [Does a tombstone in toTable() propagate to downstream groupBy/aggregate if the key never existed?](https://forum.confluent.io/t/does-a-tombstone-in-totable-propagate-to-downstream-groupby-aggregate-if-the-key-never-existed/38239)

<div class="topic-metadata">

**Author:** [@Dawid](https://forum.confluent.io/u/Dawid)\
**Replies:** 3\
**Last updated:** [2 February 2026 20:19 UTC](https://forum.confluent.io/t/does-a-tombstone-in-totable-propagate-to-downstream-groupby-aggregate-if-the-key-never-existed/38239 "2026-02-02T20:19:07Z")

</div>

Hi, I have a questions about the behavior of tombstones in Kafka Streams topology when using toTable() followed by groupBy() and aggregate(). Simplified code looks like this: \<Some 1 partition KStream\> .process(MyPro…

---

## [Issue consuming Avro](https://forum.confluent.io/t/issue-consuming-avro/38223)

<div class="topic-metadata">

**Author:** [@ernesto.costa](https://forum.confluent.io/u/ernesto.costa)\
**Replies:** 1\
**Last updated:** [14 January 2026 09:46 UTC](https://forum.confluent.io/t/issue-consuming-avro/38223 "2026-01-14T09:46:44Z")

</div>

Hi, I’m facing an issue when running a simple query trying to consume a normal topic with Avro data. The query is select \* from sdm.dk.sales.sales-orders.v1 Below you can see the show create table result. CREATE TABL…

---

## [Reading windowed topic into a global store](https://forum.confluent.io/t/reading-windowed-topic-into-a-global-store/38194)

<div class="topic-metadata">

**Author:** [@ahpo6Moh](https://forum.confluent.io/u/ahpo6Moh)\
**Replies:** 5\
**Last updated:** [10 December 2025 13:01 UTC](https://forum.confluent.io/t/reading-windowed-topic-into-a-global-store/38194 "2025-12-10T13:01:12Z")

</div>

I’ve got a consumer whose topology is defined using the Processor API. I want to read a new topic into a global store, to query that store from an existing processor. The topic is produced by another Kafka Streams applic…

---

## [Using ConsumerTimestampsInterceptor with kafka Streams](https://forum.confluent.io/t/using-consumertimestampsinterceptor-with-kafka-streams/38177)

<div class="topic-metadata">

**Author:** [@onlymohan](https://forum.confluent.io/u/onlymohan)\
**Replies:** 1\
**Last updated:** [12 November 2025 19:32 UTC](https://forum.confluent.io/t/using-consumertimestampsinterceptor-with-kafka-streams/38177 "2025-11-12T19:32:24Z")

</div>

Hi, Does using ConsumerTimestampsInterceptor work with Kafka streams? We have 2 clusters replicating data across regions with confluent replicator. If the streams is consuming from east, we want to make sure the offset…

---

## [Messages processing more than once in a exactly\_once kafka streams application](https://forum.confluent.io/t/messages-processing-more-than-once-in-a-exactly-once-kafka-streams-application/10666)

<div class="topic-metadata">

**Author:** [@Tabaldi](https://forum.confluent.io/u/Tabaldi)\
**Replies:** 7\
**Last updated:** [15 October 2025 15:59 UTC](https://forum.confluent.io/t/messages-processing-more-than-once-in-a-exactly-once-kafka-streams-application/10666 "2025-10-15T15:59:54Z")

</div>

hello everybody. I’m having problems with kafka streams recently in my application. When upgrading kafka streams from 2.8.1 to 3.4.0 i am having messages being processed more than once, even though my application has alw…

---

## [Looking for feedback? What are your most common pain points in Apache Kafka?](https://forum.confluent.io/t/looking-for-feedback-what-are-your-most-common-pain-points-in-apache-kafka/38142)

<div class="topic-metadata">

**Author:** [@szymon\_softwaremill](https://forum.confluent.io/u/szymon_softwaremill)\
**Replies:** 0\
**Last updated:** [17 October 2025 12:36 UTC](https://forum.confluent.io/t/looking-for-feedback-what-are-your-most-common-pain-points-in-apache-kafka/38142 "2025-10-17T12:36:17Z")

</div>

Hello there! :slight\_smile: I’m looking for some feedback from experienced Kafka experts like you. We’re building KafkaPilot, a tool that proactively diagnoses and resolves common issues in Apache Kafka. Currently, it c…

---

## [Kafka Streams main consumer fetch rate stays low after GlobalKTable (RocksDB) restore](https://forum.confluent.io/t/kafka-streams-main-consumer-fetch-rate-stays-low-after-globalktable-rocksdb-restore/38126)

<div class="topic-metadata">

**Author:** [@sona](https://forum.confluent.io/u/sona)\
**Replies:** 0\
**Last updated:** [1 October 2025 17:37 UTC](https://forum.confluent.io/t/kafka-streams-main-consumer-fetch-rate-stays-low-after-globalktable-rocksdb-restore/38126 "2025-10-01T17:37:27Z")

</div>

TL;DR: In Kafka Streams (3.9.1 / 4.1.0), my main stream consumer fetch rate stays low (~140k msgs/s) even after a RocksDB-backed GlobalKTable restore completes. The total fetch rate only recovers (~200k msgs/s) if I rewr…

---

## [KSQLDB Processing Time](https://forum.confluent.io/t/ksqldb-processing-time/38102)

<div class="topic-metadata">

**Author:** [@Metalmania97](https://forum.confluent.io/u/Metalmania97)\
**Replies:** 0\
**Last updated:** [13 September 2025 05:52 UTC](https://forum.confluent.io/t/ksqldb-processing-time/38102 "2025-09-13T05:52:55Z")

</div>

Hello. I am trying to test the processing delay introduced by KSQLDB and I need the processing timestamp to do this. As I understand ROWTIME is the timestamp of the kafka record. Is there a way to insert a processing ti…

---

## [Using Json Value\_format in KsqlDB I can see the topic messages null in stream](https://forum.confluent.io/t/using-json-value-format-in-ksqldb-i-can-see-the-topic-messages-null-in-stream/38096)

<div class="topic-metadata">

**Author:** [@Vidya\_Babar](https://forum.confluent.io/u/Vidya_Babar)\
**Replies:** 0\
**Last updated:** [29 August 2025 10:45 UTC](https://forum.confluent.io/t/using-json-value-format-in-ksqldb-i-can-see-the-topic-messages-null-in-stream/38096 "2025-08-29T10:45:07Z")

</div>

Using Value\_format JSON in KsqlDB I can see the topic messages null in stream using JSON Schema cloud you please tell me the reason and solution. I can see the topic messages using Value\_Format Kafka but COUNT and GRO…

---

## [Are records with a matching key from different inputs processed by the same node instances?](https://forum.confluent.io/t/are-records-with-a-matching-key-from-different-inputs-processed-by-the-same-node-instances/38087)

<div class="topic-metadata">

**Author:** [@MLaurenceFournier](https://forum.confluent.io/u/MLaurenceFournier)\
**Replies:** 1\
**Last updated:** [20 August 2025 19:41 UTC](https://forum.confluent.io/t/are-records-with-a-matching-key-from-different-inputs-processed-by-the-same-node-instances/38087 "2025-08-20T19:41:50Z")

</div>

When using “topology.addProcessor” and passing two parents with compatible record keys, are records of a given key supposed to be treated by a single processor instance? Same question with “stream.merge” followed by “str…

---

## [Problem with exactly-once semantics in kafka-streams](https://forum.confluent.io/t/problem-with-exactly-once-semantics-in-kafka-streams/38067)

<div class="topic-metadata">

**Author:** [@Mateus](https://forum.confluent.io/u/Mateus)\
**Replies:** 4\
**Last updated:** [11 August 2025 22:08 UTC](https://forum.confluent.io/t/problem-with-exactly-once-semantics-in-kafka-streams/38067 "2025-08-11T22:08:23Z")

</div>

Hello, I wanted to report an issue I encountered with my Kafka Streams application after attempting to upgrade its version from 2.8.1 to 3.4.0. My application has been running multiple Kafka Streams instances using exact…

---

## [Sliding window aggregation not counting back down to zero](https://forum.confluent.io/t/sliding-window-aggregation-not-counting-back-down-to-zero/38057)

<div class="topic-metadata">

**Author:** [@ahpo6Moh](https://forum.confluent.io/u/ahpo6Moh)\
**Replies:** 5\
**Last updated:** [6 August 2025 23:21 UTC](https://forum.confluent.io/t/sliding-window-aggregation-not-counting-back-down-to-zero/38057 "2025-08-06T23:21:27Z")

</div>

Hi there! I’m working on an aggregation using sliding windows, and I’d expect the aggregated value to return to zero after at least one window duration without event has elapsed. However, that’s not what I see in my tes…

---

## [Avoiding RecordTooLargeException on Kafka Streams](https://forum.confluent.io/t/avoiding-recordtoolargeexception-on-kafka-streams/38015)

<div class="topic-metadata">

**Author:** [@jesl](https://forum.confluent.io/u/jesl)\
**Replies:** 4\
**Last updated:** [28 July 2025 15:44 UTC](https://forum.confluent.io/t/avoiding-recordtoolargeexception-on-kafka-streams/38015 "2025-07-28T15:44:50Z")

</div>

Hello, I am using Kafka Streams to create batches of all events received on a time window. Below is the code snippet showing how I am implementing this functionality: builder.stream(inputTopic, Consumed.with(Serdes.Str…

---

## [Why Term Consume?](https://forum.confluent.io/t/why-term-consume/38007)

<div class="topic-metadata">

**Author:** [@Mannoj](https://forum.confluent.io/u/Mannoj)\
**Replies:** 1\
**Last updated:** [24 July 2025 14:02 UTC](https://forum.confluent.io/t/why-term-consume/38007 "2025-07-24T14:02:28Z")

</div>

Why Kafka has kafka-console-consumer.sh or term consumer group or consumer-app. Instead it makes more sense as kafka-console-read.sh or reader group or read-app. While data is going to stay until its TTL, it doesn’t dis…

---

## [Running a select query throwing error in control center console](https://forum.confluent.io/t/running-a-select-query-throwing-error-in-control-center-console/37970)

<div class="topic-metadata">

**Author:** [@nitin-kunal](https://forum.confluent.io/u/nitin-kunal)\
**Replies:** 0\
**Last updated:** [16 July 2025 03:31 UTC](https://forum.confluent.io/t/running-a-select-query-throwing-error-in-control-center-console/37970 "2025-07-16T03:31:09Z")

</div>

step 1 :I am using this URL to download Quick Start for Confluent Platform | Confluent Documentation step 2: got to: cd cp-all-in-one/cp-all-in-one step 3: docker compose up -d got to ksql step5: CREATE STREAM MOVE…

---

## [Kafka streams design question related to Global Ktable](https://forum.confluent.io/t/kafka-streams-design-question-related-to-global-ktable/37943)

<div class="topic-metadata">

**Author:** [@onlymohan](https://forum.confluent.io/u/onlymohan)\
**Replies:** 0\
**Last updated:** [9 July 2025 22:26 UTC](https://forum.confluent.io/t/kafka-streams-design-question-related-to-global-ktable/37943 "2025-07-09T22:26:07Z")

</div>

Current approach: We have 2 topics(A and B), we are joining data using a common field and sending to topic C. We are keeping all the data from Topic A into Global KTable, and the stream on topic B is using the topic A’…

---

## [Kafka Streams KTable Race Condition: Multiple Concurrent Updates See Same Stale State](https://forum.confluent.io/t/kafka-streams-ktable-race-condition-multiple-concurrent-updates-see-same-stale-state/37914)

<div class="topic-metadata">

**Author:** [@dongerdonger](https://forum.confluent.io/u/dongerdonger)\
**Replies:** 1\
**Last updated:** [23 June 2025 21:52 UTC](https://forum.confluent.io/t/kafka-streams-ktable-race-condition-multiple-concurrent-updates-see-same-stale-state/37914 "2025-06-23T21:52:19Z")

</div>

I’m building a conference system using Kafka Streams where users can join/leave rooms. I’m experiencing a race condition where multiple concurrent leave requests see the same stale room state, causing incorrect final re…

---

## [Multiple consumers, single partition](https://forum.confluent.io/t/multiple-consumers-single-partition/37908)

<div class="topic-metadata">

**Author:** [@GypsyCosmonaut](https://forum.confluent.io/u/GypsyCosmonaut)\
**Replies:** 1\
**Last updated:** [18 June 2025 16:54 UTC](https://forum.confluent.io/t/multiple-consumers-single-partition/37908 "2025-06-18T16:54:31Z")

</div>

We have a Kafka cluster, and in it, we have a topic named ic\_topic, which has 100 partitions. Out of these partitions, we’re seeing that a partition is assigned to multiple consumers. How is this possible ? We’re usin…

---

## [Performance impact of filtering out unnecessary records from KTable before joining](https://forum.confluent.io/t/performance-impact-of-filtering-out-unnecessary-records-from-ktable-before-joining/37840)

<div class="topic-metadata">

**Author:** [@Dawid](https://forum.confluent.io/u/Dawid)\
**Replies:** 2\
**Last updated:** [28 May 2025 09:43 UTC](https://forum.confluent.io/t/performance-impact-of-filtering-out-unnecessary-records-from-ktable-before-joining/37840 "2025-05-28T09:43:12Z")

</div>

Hello! I’m working on a Kafka Streams application that use RocksDB to save its state. I have a question regarding the performance of joins involving a KTable that may contain a large number of unnecessary records. In m…

[Next page](https://forum.confluent.io/c/stream-processing/43.md?page=1)
