Skip to content

in_kafka: Emit each record with its own tag when tag_source is record - #581

Merged
Watson1978 merged 2 commits into
fluent:masterfrom
Watson1978:in-kafka-tag-source-record
Sep 4, 2026
Merged

Watson1978 merged 2 commits into
fluent:masterfrom
Watson1978:in-kafka-tag-source-record

Conversation

@Watson1978

@Watson1978 Watson1978 commented Sep 4, 2026 •

Copy link
Copy Markdown
Contributor

in_kafka accepts tag_source record but never routes by it. consume collects the whole fetch batch into one MultiEventStream while overwriting tag in the message loop, so the single emit_stream call 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_tag already works.

def process_batch_with_record_tag(batch)
es = {}
batch.messages.each { |msg|
begin
record = @parser_proc.call(msg)
tag = record[@record_tag_key]
tag = @add_prefix + "." + tag if @add_prefix
tag = tag + "." + @add_suffix if @add_suffix
es[tag] ||= Fluent::MultiEventStream.new
case @time_source
when :kafka
record_time = Fluent::EventTime.from_time(msg.create_time)
when :now
record_time = Fluent::Engine.now
when :record
if @time_format
record_time = @time_parser.parse(record[@record_time_key].to_s)
else
record_time = record[@record_time_key]
end
else
log.fatal "BUG: invalid time_source: #{@time_source}"
end
if @kafka_message_key
record[@kafka_message_key] = msg.key
end
if @add_headers
msg.headers.each_pair { |k, v|
record[k] = v
}
end
es[tag].add(record_time, record)
rescue => e
log.warn "parser error in #{batch.topic}/#{batch.partition}", :error => e.to_s, :value => msg.value, :offset => msg.offset
log.debug_backtrace
end
}
unless es.empty?
es.each { |tag,es|
emit_events(tag, es)
}
end
end

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_stream before 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_offset unchanged and fetch the same batch on every interval.

The default tag_source topic path is unchanged and still emits a single stream.

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>
@Watson1978
Watson1978 requested a review from kenhys September 4, 2026 05:02
Comment thread lib/fluent/plugin/in_kafka.rb
Comment thread lib/fluent/plugin/in_kafka.rb Outdated
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>
@Watson1978
Watson1978 merged commit 8cc3238 into fluent:master Sep 4, 2026
61 of 63 checks passed
@Watson1978
Watson1978 deleted the in-kafka-tag-source-record branch September 4, 2026 06:08
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants