skip to content

In an ADF Mapping Data Flow, when should you change the Optimize tab away from Use current partitioning?

level: seniorimportance: nice to knowfreq 38%

answer

  1. the default is usually the right answer
  2. three tabs, and the third one matters
  3. reading a big table through one connection
  4. many files out, or exactly one
  5. round robin, hash, dynamic range, fixed range, key

basics

~20 s

Rarely. The default lets the service manage partitioning. Override it to parallelize a database source read, to co-locate a key ahead of a heavy join or aggregate, to even out badly skewed data, or to force a single output file.

solid answer

~50 s

Every transformation in a Mapping Data Flow has an **Optimize** tab offering *Use current partitioning*, *Single partition*, or *Set partitioning* with a scheme - round robin, hash, dynamic range, fixed range or key. The default carries forward whatever partitioning the stream already has and is right most of the time. Override it deliberately: - On a **database source**, source partitioning splits the read so it is not a single connection. - Before a heavy **join or aggregate**, hash or key partitioning on the grouping column co-locates matching rows. - With **skewed data and no good key**, round robin spreads rows evenly. - On a **sink**, key partitioning controls the output folder layout, and single partition is what forces one output file. Single partition funnels everything through one worker, so use it only where output shape demands it, and measure in the monitoring detail before and after.

code

text · 9 lines
text
Source SqlOrders            Optimize: Set partitioning
                              Scheme  : Dynamic range
                              Column  : order_id

Join OrdersToCustomers      Optimize: Use current partitioning

Sink CuratedOrders          Optimize: Set partitioning
                              Scheme  : Key
                              Column  : order_date   -- one folder per date

go deeper

for a junior

Know that an Optimize tab exists on each transformation and that the default handles partitioning for you. Do not change it without a reason you can articulate.

for a middle

Name the schemes - round robin, hash, dynamic range, fixed range, key - and explain what each does to how rows are distributed, plus why single partition is the setting behind single-file output.

for a senior

Show the diagnostic loop: read per-stage rows and timings in monitoring, form a hypothesis such as an unparallelised source read or a skewed join key, change one setting, and re-measure on comparable data.

for a principal

Decide when tuning stops being worthwhile. If a flow needs careful hand partitioning to meet its window, that is evidence the transformation may belong in the warehouse or on compute the team controls.

## What the Optimize tab offers Every transformation node in a Mapping Data Flow - source, join, aggregate, sink and everything in between - has three tabs, and the third is **Optimize**. It offers three choices: - **Use current partitioning** (the default): keep whatever partitioning the incoming stream already has and let the service decide. - **Single partition**: collapse the stream into one partition. - **Set partitioning**: choose a scheme explicitly - *round robin* (with a partition count), *hash* (on chosen columns), *dynamic range*, *fixed range*, or *key* (a column whose distinct values define the partitions). The important cultural point first: the default is the right answer for the overwhelming majority of transformations. Setting partitioning on every node because you once read that partitioning matters is a reliable way to make a flow slower. The Optimize tab is a targeted instrument, not a checklist item. ## The cases that justify an override **Reading a relational source in parallel.** A database source read that is left to a single stream can be the bottleneck of the whole flow. Source-side partitioning splits the read into ranges or hashes so several readers pull concurrently. This needs a column with reasonable distribution - a monotonic key or a well-spread numeric or date column. Partitioning on a column with a handful of distinct values just creates a few enormous partitions. **Co-locating data ahead of an expensive operation.** Joins and aggregates work on matching keys. If the stream is already partitioned by the join or grouping key, matching rows are together and the operation does less movement. Hash or key partitioning on that column before a heavy join can pay for itself. This is also the case where the choice is easiest to get wrong: if the key is skewed, you have simply moved the problem into one very large partition. **Fighting skew with no usable key.** When no column distributes well, round robin with an explicit partition count spreads rows evenly by construction. It destroys locality, so it is the wrong choice immediately before a join on a key, and a reasonable one before independent per-row work. **Controlling the shape of the output.** This is the most concrete and most frequently asked case. A sink writing to a data lake produces files per partition, so partitioning drives the output layout: key partitioning by, say, a date column gives you a folder or file per date, while a stream with many partitions gives you many files. When the requirement is a single output file - a downstream consumer that cannot glob a folder, or a small extract for a partner - the sink's Optimize tab must be set to single partition, which is also what the sink's single-file naming option depends on. ## The cost of single partition Single partition is the setting people reach for casually and regret. Everything funnels through one worker: one partition's worth of memory, no parallelism, and a hard ceiling on throughput. On a small result set that is fine and is exactly what you want for a single-file extract. On a large one it can turn a comfortable flow into a memory failure. If a downstream system merely dislikes many *tiny* files, the better answer is usually to reduce partition count moderately or to partition by a meaningful key, not to collapse to one. ## Measure, then change The data flow monitoring detail for a run shows the transformation graph with row counts and per-stage timings. That is where you find out whether a stage is genuinely slow or whether the activity duration was dominated by cluster acquisition - a distinction that decides whether partitioning is even the right lever. Change one setting, rerun on comparable data, and compare the same stage. Partitioning changes interact, so changing several nodes at once tells you nothing about which one helped. It is also worth remembering what partitioning does not fix. It does not reduce cluster startup time. It does not compensate for reading far more data than you need - a filter pushed into the source query beats any partitioning scheme. And it does not rescue a flow whose real problem is that a fundamentally SQL-shaped transformation is being done on an external engine at all. ## How to answer this in an interview Say that the default is correct by default, name the schemes, then give two or three concrete triggers - parallelising a database read, co-locating a join key, and shaping sink output including the single-file case - and finish with the caution about single partition and the habit of measuring in the monitoring view. That reads as someone who has tuned a real flow rather than someone reciting options.

  • A sink is producing thousands of tiny files in the data lake. Is single partition the fix?
    Usually not. Single partition removes all write parallelism and can exhaust one worker's memory on a large result. Prefer reducing the partition count or partitioning by a meaningful key such as load date, so you get a sensible number of right-sized files. Reserve single partition for genuinely small outputs that must be one file.
  • You hash-partition on the join key and the stage gets slower. What happened?
    Most likely skew: if a few key values dominate, hashing sends all of their rows to the same partition, so one task does nearly all the work while the rest idle. Check per-stage row counts in the monitoring detail. Either fall back to the default, filter or handle the hot values separately, or spread rows with round robin when locality is not needed.
  • Where do you look to tell a slow transformation from a slow activity?
    The data flow monitoring detail for the run, which shows each transformation with rows processed and time spent. If the stages account for a small share of the activity duration, the rest went to acquiring compute, and partitioning is the wrong lever entirely.

saying these in an interview costs you the question

  • Sets explicit partitioning on every transformation by habit
  • Uses single partition to reduce file count on large outputs
  • Hash-partitions on a skewed key and expects balance
  • Thinks partitioning reduces cluster startup time
  • Changes several Optimize settings at once and cannot attribute the result

context