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)

Reply via email to