Skip to content
Back to skills

Stream Processing Runtime Performance

ASecurity

Operating Kafka Streams and Apache Flink for predictable throughput, state and recovery: separating their execution models, sizing partitions or operator parallelism, diagnosing backpressure, bounding native state, and relating commits or checkpoints to result visibility. Use when one partition or operator limits a pipeline, checkpoints grow or stall, RocksDB drives RSS outside the heap, exactly-once changes latency, or effective runtime configuration differs from declared settings. Generic t...

  • 2 stars
  • 0 votes
  • 0 copies
  • 2 views
  • Added September 19, 2026
developmentgojavaawsapiperformance

Works with

  • api

Security analysis

A100/100

Pro scans all 4 files and shows the line behind each finding

Scanned September 29, 2026

npx -y skills add robsonkades/agent-skills --skill stream-processing-runtime-performance --agent claude-code

Installs into .claude/skills of the current project.

Are you the author of Stream Processing Runtime Performance?

Add the live security badge to your README. It updates with every re-scan.

Security grade badge for Stream Processing Runtime Performance
[![Security: A — Skills Directory](https://www.skillsdirectory.com/api/skills/robsonkades-stream-processing-runtime-performance/badge)](https://www.skillsdirectory.com/skills/robsonkades-stream-processing-runtime-performance)

More formats (shields.io, HTML) on the badges page. Keep it an A: scan every change in CI with Pro.

Download with Pro
SKILL.md
---
name: stream-processing-runtime-performance
description: >
  Operating Kafka Streams and Apache Flink for predictable throughput, state and recovery:
  separating their execution models, sizing partitions or operator parallelism, diagnosing
  backpressure, bounding native state, and relating commits or checkpoints to result visibility.
  Use when one partition or operator limits a pipeline, checkpoints grow or stall, RocksDB drives
  RSS outside the heap, exactly-once changes latency, or effective runtime configuration differs
  from declared settings. Generic topology and event-time semantics belong to
  streaming-pipeline-topologies; plain Kafka consumer loops to kafka-consumers-in-java.
---

# Stream-Processing Runtime Performance

## Purpose

Turn a streaming symptom into the runtime resource, state or recovery mechanism that owns it.
Kafka Streams is an embedded partition-to-task library; Flink is a distributed operator graph.
They share concepts but not a tuning surface.

## Common contract

Use the parts of this inventory relevant to the question, reusing adequate supplied evidence.
A task-count, TTL or configuration API explanation need not trigger unrelated captures or a change.
Record relevant engine and connector versions, topology/job graph, input partitions and key distribution,
operator/task parallelism, offered/completed events and bytes, per-partition age/lag, backpressure,
state location/size, allocation rate, heap/RSS/container limits, checkpoint/commit configuration,
processing guarantee, sink visibility and recovery objectives.
Inspect the project's JDK/toolchain and engine/connector support matrix; the references use Kafka
4.0 and Flink 2.0 as evidence baselines, not mandatory upgrades. Preserve topology IDs, partitioning,
state schemas and deployment authority when changing runtime settings.

## Workflow

Apply the steps needed to establish the requested diagnosis or decision. Missing evidence leaves
the affected mechanism unresolved; it does not invalidate independent documented facts.

1. Draw source partitions through every shuffle/operator to state and sinks. Mark ownership,
   serialization and atomic recovery boundaries.
2. Locate the limiting partition or operator with aligned rates, waits and resource evidence; aggregates hide skew and head-of-line
   blocking.
3. Separate durable backlog from runtime backpressure. Kafka lag can grow without slowing producers;
   Flink operator credit/backpressure propagates within the job graph.
4. For a capability/default question, use the exact version's API, configuration definition or
   source. For deployed behavior, reconcile declared and effective settings, component ownership
   and precedence using relevant runtime evidence. Verify the exact key/parser path before claiming
   an unknown or deprecated setting is ignored, rejected, warned about or mapped to another key.
5. Account for heap allocation and native/file-backed state separately inside the same container.
6. Retain adequate behavior or propose a bounded, evidenced change. Reuse relevant validation or
   compare the affected steady-state, failure, restore or rescale contract in an authorized scope;
   a narrow answer does not require every campaign. Do not report a proposed test as executed.

## Recovery completion

Define whether recovery means task availability or restored output timeliness. For the latter,
observe restart/assignment, state restoration, backlog replay and sink visibility; phases may
overlap, so do not blindly sum their durations. A running task or restored store alone does not
establish the output-age objective. Include required replay and arrivals during downtime in the
backlog at resumption, without counting the same work twice.

For a stable work mix at one bottleneck, a backlog `B` drains in approximately
`B / (mu - lambda)` when sustained useful processing capacity `mu` exceeds ongoing arrival rate
`lambda`. Use the same work units and stage boundary for all three; restore bytes/second and
post-aggregation output counts are not interchangeable with source events/second. Measure capacity
under relevant recovery conditions. If `mu <= lambda`, waiting cannot drain a positive backlog
under those assumptions; a shorter state restore alone does not demonstrate processing headroom.
Check the limiting partition/operator, since aggregate spare capacity may not be assignable there.

Validate backlog age and useful sink results while ingress continues, against the requested recovery
objective and source-retention limits. Reuse adequate existing recovery evidence. For broader capacity
or admission changes, pass the stage rates, skew, resource limits and recovery target to
`capacity-planning`; if unavailable, report the shortfall and a scoped measurement/change to validate.

## Rules

- Size Kafka partitions from representative workload and per-partition capacity for each required
  consumer group, with producer, broker/network/storage, key/order and recovery constraints plus
  growth/failure headroom. Powers of two or broker multiples are placement heuristics, not laws.
- Size Flink per operator. Raising global parallelism cannot repair one serialized sink, skewed key
  or blocking call.
- Exactly-once has a declared boundary. External effects outside the engine transaction/checkpoint
  need idempotency, fencing or their own transaction protocol.
- Checkpoint success is not automatically sink visibility. State when a two-phase sink commits and
  include that delay in the output-latency contract.
- State size does not determine Java heap by itself. Allocation, retained live state and object
  lifetimes affect GC; RocksDB caches/write buffers and file pages have distinct native/resident
  accounting. Managed memory is a budget, not an additional RSS category to sum twice.
- Backpressure is a symptom location, not necessarily the root. Trace downstream as a starting
  hypothesis, then correlate source starvation, shared CPU/network/storage, async work and
  checkpoint/sink coupling. One low-output vertex does not prove an independent local bottleneck.
- Busy task time is not CPU utilization: blocking callbacks and chained operators need stack,
  I/O and downstream evidence. Healthy throughput with stalled watermarks can instead be an
  event-time/idleness problem; do not add parallelism solely because a window emits nothing.
- Return the scoped answer or retained/proposed decision with its evidence and uncertainty. For
  diagnosis or change, include the competing cause, relevant effective settings and validation;
  distinguish successful visible results from input consumption, dropped events and retries.

## References

- [Kafka Streams operation](references/kafka-streams.md) — read for task/partition capacity,
  transactions, state stores, standby replicas or commit visibility.
- [Flink operation](references/flink.md) — read for operator backpressure, checkpoints, RocksDB,
  watermarks, savepoints or connector lifecycle.

Files in this skill

  • SKILL.md5.4 KB
  • references/flink.md5.2 KB
  • references/kafka-streams.md3.2 KB
  • skill.yaml1.6 KB

Attribution

Is this your skill, or is something wrong with this listing? Request removal or report an issue. Author removals are honored within 72 hours.

Comments

Loading comments…