sanghyeok An created KAFKA-21171:
------------------------------------
Summary: Repeated task offset snapshot rebuilding adds overhead
with the Streams rebalance protocol
Key: KAFKA-21171
URL: https://issues.apache.org/jira/browse/KAFKA-21171
Project: Kafka
Issue Type: Improvement
Reporter: sanghyeok An
Assignee: sanghyeok An
h2. Summary
With {{{}group.protocol=streams{}}}, {{StreamThread}} rebuilds the task offset
sum snapshot on every processing loop iteration, even when the offsets and task
states have not changed. This performs task enumeration, map copying,
filtering, and allocation independently of whether the next heartbeat will
report new offsets.
A local JMH benchmark shows that the cost of one refresh grows approximately
linearly with task count. In a separate 10-minute application profile, this
path accounted for *580 of 1,207 Java execution samples (48.1%)* from the
Streams processing thread, including descendant calls.
h2. Current behavior
{{TaskManager.maybeUpdateTaskOffsetSumSnapshot()}} performs the following work:
# Copies the process-wide local offset map from {{{}StateDirectory{}}}.
# Builds a map of this thread's tasks through {{{}allTasks(){}}}, including
tasks held by the state updater and tasks awaiting initialization.
# Removes running active tasks and overlays live changelog offset sums for
other eligible tasks.
# Builds a separate map with heartbeat task IDs and publishes a read-only
snapshot.
The snapshot is refreshed on each loop iteration in both processing-loop
variants. Heartbeat code separately checks assignment changes, the offset
reporting interval, warm-up progress, and whether the values changed. These
checks limit offset transmission but do not avoid the preceding snapshot
construction.
The copied local offset map is in memory; this path does not read task state
from disk on every iteration. Under {{{}group.protocol=classic{}}}, this
snapshot refresh method returns immediately; classic reports offset sums
through its subscription metadata path.
h2. Environment
* Kafka {{{}4.5.0-SNAPSHOT{}}}, checkout
{{{}e0b738a78819d95dd7779ec3d98effbb9f6a0f9c{}}}.
* macOS AArch64; JFR reports 14 hardware threads.
* JMH: JMH 1.37, JDK 25.0.3.
* Application profile: JDK 21.0.5+11-LTS.
* Affected released versions have not been established.
h2. JMH results
{{TaskOffsetSumSnapshotBenchmark.refreshSnapshot}} measures one refresh with
unchanged offsets. All tasks belong to one stream thread and each has one
changelog partition.
The fixture uses the real task manager, task registry, state manager, state
updater, and local offset map. It does not start a broker or perform disk I/O.
RUNNING tasks are registered as active tasks and produce an empty final
snapshot. RESTORING tasks remain in the state updater input queue; the updater
thread is not started.
Settings: one worker, three forks, {{{}-Xms2g -Xmx2g{}}}, five 2-second warmup
iterations and five 2-second measurement iterations per fork, {{{}avgt{}}}, and
{{{}-prof gc{}}}. There are 15 measurements per parameter combination. Time
errors below are the JMH 99.9% confidence interval half-widths.
||Task count||Task state||Time (us/op, mean ± error)||Allocation (B/op)||
|0|RUNNING|0.257 ± 0.002|2,120.001|
|0|RESTORING|0.259 ± 0.001|2,120.001|
|100|RUNNING|8.481 ± 0.057|31,890.708|
|100|RESTORING|22.821 ± 0.111|77,048.131|
|1,000|RUNNING|89.110 ± 1.960|346,829.769|
|1,000|RESTORING|248.191 ± 1.702|831,546.901|
|5,000|RUNNING|429.778 ± 5.789|1,732,682.127|
|5,000|RESTORING|1,304.563 ± 22.457|4,111,075.094|
The RUNNING cases show that copying and filtering still incur a cost when all
local entries are excluded from the final snapshot. Allocation is measured per
refresh, not retained heap size.
The results were collected in two runs: task count 1,000, and task counts
0/100/5,000. For example, the latter used:
{code:bash}
java -jar jmh-benchmarks/build/libs/kafka-jmh-benchmarks-4.5.0-SNAPSHOT-all.jar
\
TaskOffsetSumSnapshotBenchmark \
-p taskCount=0,100,5000 -p taskState=RUNNING,RESTORING \
-f 3 -wi 5 -i 5 -w 2s -r 2s \
-prof gc -foe true \
-rf json -rff jmh-benchmarks/build/task-offset-snapshot-scaling.json
{code}
Benchmark source: [^TaskOffsetSumSnapshotBenchmark.java]
Raw results: [^task-offset-snapshot-1000.json],
[^task-offset-snapshot-scaling.json]
h2. Application profile and flame graph
The application uses one sub-topology:
{noformat}
snapshot-bench-input
-> groupByKey()
-> count(persistent key-value store)
{noformat}
The configured test uses {{{}group.protocol=streams{}}},
{{{}num.stream.threads=1{}}}, and {{{}poll.ms=10{}}}. JFR was recorded with
{{{}settings=profile{}}}, a 60-second start delay, and a 600-second duration.
The intended workload was 100 tasks with no continuous input; the exact
assigned task count and input activity during this recording were not
independently captured.
The recording contains 1,225 {{jdk.ExecutionSample}} events, including 1,207
from {{{}StreamThread-1{}}}. The flame graph uses these Java execution events
only. Native-method, allocation, and waiting events are excluded; none of the
recorded execution stacks are marked truncated.
||Processing-thread call path, including descendants||Samples||Share of 1,207
samples||
|TaskManager.maybeUpdateTaskOffsetSumSnapshot|580|48.1%|
|TaskManager.process|255|21.1%|
|StreamThread.maybeCommit|157|13.0%|
|TaskManager.recordTaskProcessRatio|128|10.6%|
|TaskManager.punctuate|39|3.2%|
|pollPhase and other paths|48|4.0%|
The {{process}} and {{maybeCommit}} paths include readiness checks; their
presence does not establish that records were processed or commits completed.
Within the 580 snapshot-path samples, the sampled call sites were:
||Snapshot operation||Samples||
|allTasks() and task iteration preparation|284|
|Copying the StateDirectory offset map|194|
|Removing running active tasks from the copied map|98|
|Other|4|
Descendant frames include {{{}TaskId.hashCode(){}}}, {{HashMap}}
insertion/removal, {{ConcurrentHashMap}} traversal, and state-updater task
enumeration.
*48.1% is a share of Java execution samples, not whole-process CPU utilization
or the fraction of a CPU core consumed.*
Separately, cumulative thread CPU counters in the first and last recorded
thread dumps increased by 63.956 seconds over 543.632 seconds of elapsed time
for the Streams processing thread. This corresponds to approximately *11.8% of
one CPU core* over that interval. The sample share does not establish an exact
CPU time or potential saving for the snapshot method.
Flame graphs: [^stream-thread.svg], [^all-java.svg]. Download and open the SVG
in a browser to zoom and search for {{{}maybeUpdateTaskOffsetSumSnapshot{}}}.
h2. Interpretation and limits
The microbenchmark demonstrates repeated work with unchanged inputs, and the
application profile identifies snapshot construction as a prominent Java
execution path in the measured run.
An earlier 120-second recording with {{poll.ms=100}} contained only 27 Java
execution samples and no snapshot-path samples. It is too sparse to establish
the snapshot's CPU share. The application ID also changed between runs, and
task/input conditions were not fully recorded, so these profiles are not a
controlled comparison of poll intervals.
The JMH and application runs used different JDK versions. No production
throughput regression, lock contention bottleneck, or measured improvement is
established by these results.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)