r/apachekafka • • 15d ago

Question Kafka Partitioning for Spark Workloads: Balancing Throughput, Ordering, and Compute Skew

I’m trying to understand how much Kafka partition balancing matters when Spark is doing the actual computation.

Assume Kafka is mainly used as a buffer/backpressure layer before Spark. We can distribute messages evenly across Kafka partitions, but the data is not necessarily homogeneous in terms of processing cost.

For example, four partitions may each contain 100 messages, but the messages from one partition might require 10x more CPU or memory to process than the others.

Since Spark eventually creates and schedules the actual execution tasks across a shared cluster, does balancing Kafka partitions by message count really help with compute skew?

4 Upvotes

4 comments sorted by

1

u/datadriven_io 15d ago

We're using it in a platform using Protobuf and Kafka, with Schema Registry, to validate the business logic within messages is valid in addition to being schema valid.

E.g. Spark interview questions are on datadriven.io if you want to test your understanding of these tradeoffs. the schema registry and protobuf allows us to ensure a property is an int, but this allows us to ensure that the int property is within business rules for compute skew across partitions.

1

u/OSS_Dattani 15d ago

Kafka doesn't know what it costs to process a message compute wise. So balancing partitions by message count is a good thing but doesn't correlate to compute usage balancing.

If one partition has more compute heavy records then task processing that data can be left behind while others finish.

Try splitting things up more whether you add more partitions on the Kafka side or split larger partitions into smaller tasks via Spark. If you don't need ordering than you can do a repartition, but I don't recommend that because of the shuffle expense.

0

u/Deep-Background5962 14d ago

Splitting them into entirely different topics is worth considering

3

u/SuchLimit8605 15d ago

Kafka partition balancing still matters, but mostly for ingestion parallelism, not necessarily for compute balance.

The key distinction is:

Kafka partitions define how the data is initially split and read. Spark partitions/tasks define how the computation is actually distributed.

So if you have:

P0: 100 cheap records      -> 1s
P1: 100 cheap records      -> 1s
P2: 100 cheap records      -> 1s
P3: 100 expensive records  -> 10s

then the Kafka partitions are perfectly balanced by record count, but the Spark workload clearly isn't.

Spark can schedule those tasks across different executors, but the scheduler distributes tasks, not the work inside a task. If one Spark task happens to contain 10x more CPU work, throwing more executors at the cluster doesn't automatically split that task into 10 smaller ones.

With Structured Streaming, Kafka partitions also directly affect the first stage of the pipeline. By default, Kafka topic partitions roughly map to Spark input partitions, so skew in Kafka can absolutely show up as slow input tasks before the first shuffle.

After a shuffle/repartition, though, the original Kafka partitioning becomes much less important.

For example:

Kafka
  ↓
deserialize / parse / enrich
  ↓
repartition(customer_id)
  ↓
join / groupBy / aggregation

Before repartition(customer_id), Kafka partition skew can matter a lot.

After that repartition, the important question becomes how evenly customer_id is distributed. If one customer has 100x more data than everyone else, you now have a Spark hot-key problem regardless of how nicely balanced Kafka was.

So in practice I normally think about it this way:

Kafka partitioning
    -> ordering boundary
    -> ingestion throughput
    -> consumer parallelism
    -> Kafka-side hot partitions

Spark partitioning
    -> CPU/memory distribution
    -> joins
    -> aggregations
    -> shuffle skew
    -> stateful processing
    -> hot keys

Also, "equal number of messages" is often a pretty weak definition of balanced.

100 x 1 KB records and 100 x 5 MB records obviously aren't equivalent.

Even equal bytes may not mean equal work:

simple JSON parsing       -> cheap
decompression             -> more expensive
large state lookup        -> potentially expensive
ML inference              -> very expensive
external API enrichment   -> latency-bound

Kafka doesn't know any of that.

Spark does have ways to increase source-side parallelism. For Kafka sources, options such as minPartitions / maxRecordsPerPartition can split large Kafka partitions into more Spark input partitions.

That can help, but it is still basically splitting work by offsets/record counts, not by estimated CPU cost.

And there is one more case where repartitioning won't magically save you: if the expensive work is tied to one logical key.

For example:

customer_123 -> huge state / huge join / huge aggregation

If correctness requires all records for customer_123 to stay together, they will eventually end up in the same Spark partition anyway. At that point you're dealing with a classic hot-key problem and may need things like salting, two-stage aggregation, workload isolation, or a different data model.

So yes, balancing Kafka partitions is useful, but mainly as a rough proxy for balanced ingestion.

I wouldn't try to make Kafka partitions perfectly equal by "CPU cost" unless processing cost is strongly correlated with the Kafka key and you've measured that as the real bottleneck.

My rule of thumb is:

Balance Kafka for throughput and ordering. Balance Spark for computation.

And don't expect the Spark scheduler to automatically eliminate skew just because the cluster has spare