in_kafka: Emit each record with its own tag when tag_source is record - #581
Merged
Merged
Conversation
consume built a single MultiEventStream and overwrote tag inside the message loop, so emit_stream ran once with the tag of the last parsed record. Every record of the same fetch batch was routed to that tag, which means tag_source record never worked in in_kafka. Group the records by tag and emit one stream per tag, the same way in_kafka_group#process_batch_with_record_tag does. The tag_source topic path keeps emitting a single stream. Skip a record whose tag field is not a String. Emitting per tag makes one such record raise out of emit_stream before the remaining streams are emitted, which leaves the offset unchanged and re-emits the earlier tags on every interval. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Shizuo Fujita <fujita@clear-code.com>
kenhys
requested changes
Sep 4, 2026
kenhys
reviewed
Sep 4, 2026
emit and the offset update shared one guard, so a batch whose records all failed to parse or carried an invalid tag left @next_offset where it was. consume fetches from @next_offset, so the same batch came back on every interval and the same warnings were logged forever. Reject an empty tag as well. It reaches no match pattern, so the record is dropped by the router with only a "no patterns matched" warning. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Shizuo Fujita <fujita@clear-code.com>
kenhys
approved these changes
Sep 4, 2026
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
in_kafkaacceptstag_source recordbut never routes by it.consumecollects the whole fetch batch into oneMultiEventStreamwhile overwritingtagin the message loop, so the singleemit_streamcall at the end uses the tag of the last parsed record. Every record of the batch is emitted with that tag, and records that carry a different tag field are silently routed to the wrong destination.This groups the records by tag and emits one stream per tag, the same way
in_kafka_group#process_batch_with_record_tagalready works.fluent-plugin-kafka/lib/fluent/plugin/in_kafka_group.rb
Lines 267 to 310 in e6be8f1
Records whose tag field is not a String, or is an empty string, are now skipped with a warning. Emitting per tag makes one such record raise out of
emit_streambefore the remaining streams are emitted, so the offset stays where it is and the earlier tags are emitted again on every interval. The offset now also advances when every message of a batch is skipped, which used to leave@next_offsetunchanged and fetch the same batch on every interval.The default
tag_source topicpath is unchanged and still emits a single stream.