skip to content

How does MirrorSourceConnector achieve parallelism, and how do tasks.max and source partitions affect replication throughput?

level: seniorimportance: should knowfreq 40%

answer

  1. Unit of work = source topic-partition
  2. Parallelism = min(tasks.max, #partitions)
  3. One partition -> one task -> order preserved
  4. Workers spread tasks; tasks.max caps count
  5. Tune producer batch/compression for WAN

basics

~20 s

MirrorSourceConnector divides the source topic-partitions it must replicate across tasks. Parallelism is bounded by tasks.max and by the number of source partitions — you never get more useful tasks than partitions, since each partition is handled by one task.

solid answer

~40 s

MirrorSourceConnector enumerates all source topic-partitions matching the filters, then distributes them round-robin across up to `tasks.max` tasks; each task runs a consumer+producer pipeline mirroring its assigned partitions. So effective parallelism = min(tasks.max, total source partitions). Setting tasks.max above the partition count wastes slots. Throughput also depends on the underlying producer/consumer configs MM2 passes through (`<flow>.producer.*`, `<flow>.consumer.*` — batching, compression, `max.poll.records`), and on Connect worker count, since tasks are balanced across workers by the Connect rebalance protocol. Because each task owns whole partitions, partition-level ordering is preserved end-to-end. Checkpoint and heartbeat connectors have far lighter task needs. To scale, you increase tasks.max, add Connect workers, and ensure source topics have enough partitions; you also tune producer batching/compression to saturate the inter-cluster link.

go deeper

for a junior

Know MM2 runs multiple tasks and that tasks.max sets an upper bound on parallelism.

for a middle

Explain that one partition maps to one task and parallelism is min(tasks.max, partitions).

for a senior

Discuss worker vs task scaling, producer/consumer tuning for WAN, and straggler partitions.

for a principal

Capacity-plan an MM2 deployment: partition counts, tasks.max, worker fleet sizing, and producer tuning to saturate cross-DC links while preserving ordering.

## The parallelism model Kafka Connect connectors scale by splitting work into **tasks**, and tasks are spread across **workers** (JVMs) by Connect's group-membership/rebalance protocol. For MirrorSourceConnector, the unit of work is a **source topic-partition**. The connector lists every topic-partition that passes the topic allow/deny filters, then assigns those partitions across up to `tasks.max` tasks (Connect's `taskConfigs(maxTasks)` returns the partition groupings). Each task opens a consumer on the source cluster for its partitions and a producer to the target cluster, forming an independent mirroring pipeline. ## The fundamental bound Because one partition is mirrored by exactly one task (to preserve per-partition ordering), the **maximum useful parallelism is min(tasks.max, number_of_source_partitions)**. - If you have 12 source partitions and set `tasks.max=50`, you get at most 12 active tasks; the rest are idle. - Conversely if `tasks.max=2` over 12 partitions, each task handles 6 partitions serially within its poll loop. So the two levers are: raise `tasks.max` and ensure the **source topics actually have enough partitions**. ## Where tasks run `tasks.max` is per-connector; the actual tasks are balanced across all Connect workers in the (dedicated or shared) cluster. Adding workers gives the existing tasks more CPU/network headroom and lets Connect spread them out, but it does **not** create more tasks than `tasks.max` allows. So horizontal scaling is two-dimensional: - `tasks.max` (how many parallel pipelines) and - worker count (where they run). ## Throughput tuning beyond task count MM2 lets you pass producer/consumer configs per flow with prefixes like `<src>-><tgt>.producer.override.*` (or the cluster-level `<alias>.producer.*`). Key knobs: - `batch.size`, `linger.ms`, `compression.type` (often `lz4`/`zstd` to shrink cross-DC bytes), `buffer.memory`, - and on the consumer side `max.poll.records`, `fetch.max.bytes`, `receive.buffer.bytes`. Because replication frequently crosses a high-latency WAN, large batches + compression + adequate in-flight requests matter more than raw task count once partitions are covered. ## Ordering and correctness Since a partition is never split across tasks, **per-partition order is preserved** through the mirror. (Order is not preserved *across* partitions, same as normal Kafka.) Records keep their original partition by default (MM2's default `IdentityReplicationPolicy`-style partitioning maps partition i->i), which is what lets offset translation work meaningfully. ## The other connectors MirrorCheckpointConnector and MirrorHeartbeatConnector are light: heartbeats are a single low-rate stream; checkpoints scale with the number of consumer groups, not data volume. They rarely need high `tasks.max`. ## Common mistakes / edge cases - (1) Cranking `tasks.max` without adding partitions or workers yields no gain and can cause excessive rebalances. - (2) A few very-hot partitions create stragglers because a task can't subdivide a single partition — the fix is more source partitions, not more tasks. - (3) Forgetting to tune producer batching leaves a high-latency WAN underutilized even with many tasks. - (4) Frequent connector restarts trigger Connect rebalances that briefly pause all tasks.

  • You set tasks.max=20 but throughput doesn't improve past 8 tasks. What's the likely cause?
    The source topics being mirrored have only ~8 total partitions, so MM2 can't create more than ~8 useful tasks (one partition per task). Add partitions on the source, or accept the bound. Also verify Connect has enough workers to actually run them.
  • How is per-partition ordering preserved across the mirror?
    Each source partition is assigned to exactly one task, which mirrors it in order to the corresponding target partition (default i->i mapping). Since a partition is never split across tasks, original ordering within that partition is maintained.

saying these in an interview costs you the question

  • Claiming you can scale a single hot partition by adding tasks — a partition is owned by one task and can't be subdivided.
  • Saying more Connect workers create more tasks — workers only host tasks; tasks.max caps the count.
  • Ignoring that effective parallelism is bounded by source partition count.
  • Thinking task count alone determines WAN throughput, ignoring producer batching/compression tuning.

context