skip to content

Why does handing records to a user-supplied function cost more than that function's own work, and what makes a cross-language one worse?

level: seniorimportance: should knowfreq 44%

answer

  1. two representations, one boundary
  2. decode in, encode out, per record
  3. cost scales with records, not logic
  4. a different language means another process
  5. batch, narrow, or do not cross

basics

~20 s

Where the engine holds records packed, host-language code cannot read them: each record is decoded into objects going in and encoded back coming out. That cost is per record, so a trivial function can be dominated by it. A different language adds a process hop and a second conversion.

solid answer

~50 s

Where an engine keeps **engine-managed binary records** - a packed byte layout it allocates and interprets itself - its own operators read fields at offsets without ever building an object. A user-supplied function written in the worker's language cannot do that, so the engine decodes each record into **host-runtime objects**, runs the function, and encodes the result back. Two conversions per record, plus the allocations the decode performs. The cost scales with record count and field count, not with how much the function computes, which is why a function doing almost nothing can dominate a step. A function in a *different* language is worse again: it runs in a second process, so records are copied across a process boundary and converted into something both sides read, and that process's memory is real footprint the engine's accounting never counted. On an engine that holds ordinary objects throughout, there is no conversion to pay and the cost lies elsewhere.

go deeper

for a junior

Know that a function you write is not always free to receive a record: on some engines the record has to be converted into an object first, and converted back afterwards.

for a middle

Explain that the conversion is charged per record and per field rather than per unit of computation, and that a cross-language function additionally runs in a separate process with a copy in each direction.

for a senior

Demonstrate the diagnosis - identity-function comparison, records-in against duration - and the ordering judgment: filter and project before the crossing, batch it, fuse consecutive crossings, or avoid it by expressing the work in engine operators.

for a principal

Decide where crossings are permitted at all. Allowing per-record functions in another language is a capability authors will use everywhere, and it imports both a conversion bill and memory your engine's accounting cannot see.

## The boundary Two representations sit on either side of this question. **Engine-managed binary records** are records the engine keeps as a packed byte layout it allocates and interprets itself, so a field is read at a known offset. **Host-runtime objects** are records held as ordinary objects of whatever language the **worker** - one process on one machine running some of the job's work - happens to run. An engine's own operators, where it has a packed form, can filter, compare, hash and copy records by reading bytes. A function you supply cannot: it expects an object with fields you can name. So at the point where the pipeline enters your function, the engine must **decode** the packed bytes into objects, and at the point where it leaves, **encode** the objects back into the packed layout for whatever comes next. Two conversions, per record, per crossing. ## Why the cost tracks records, not work The conversion is proportional to the number of records and the number of fields materialised, and it is indifferent to what the function does with them. That produces the counter-intuitive result people actually meet: - A function that does nothing but return a field unchanged can still make its step the most expensive in the job. - Making the function's own logic faster changes almost nothing, because the function was never the cost. - The same function over a record of thirty fields costs far more than over a record of three, even if it reads one field either way, on engines that materialise the whole record. The measurement that settles it is to run the same step with the function replaced by one that returns its input untouched. Whatever time remains above the surrounding steps is the boundary, not the logic. ## The cross-language case When the function is written in a language other than the one the workers run, it cannot execute inside the worker process at all. It runs in a second process on the same machine, and the records have to get there and back. | Cost | Same-language function | Cross-language function | |---|---|---| | Conversion | Decode to objects, encode back | Encode into a form the other runtime reads, then decode there - and the same in reverse | | Movement | None; same process memory | A copy across a process boundary in each direction | | Memory accounting | Inside the worker's budget | The second process's footprint is real, and the engine's accounting never counted it | | Per-call overhead | A function call | An inter-process round trip, unless calls are batched | The memory row is the one that catches people out. The second process's footprint still counts against the ceiling the platform enforces on the whole process - a ceiling applied by the platform rather than the engine, enforced by killing the process outright rather than by failing one allocation. So a worker that the engine believes is comfortably inside its budget can be killed because of memory the engine never accounted for. How a worker's budget is divided, and the gap between the engine's accounting and the process's real footprint, is a subject of its own; what belongs here is that a cross-language boundary manufactures that gap. ## Reducing it 1. **Do not cross.** If the work can be expressed with the engine's own operators, no conversion happens at all. This is by far the largest win and it is usually available for filtering, projection, arithmetic and standard aggregation. 2. **Cross in batches.** A surface that hands the function many records at once amortises per-call overhead and lets the conversion be done in bulk rather than per record. 3. **Share a layout.** Where both sides can read the same in-memory batch layout, one encode-and-decode pair disappears and the transfer becomes a bulk copy rather than a per-record translation. 4. **Cross as late and as narrow as possible.** Filter and drop fields before the function so that fewer records with fewer fields are converted. Ordering the pipeline this way often beats optimising the function. 5. **Cross once.** Three consecutive per-record functions can mean three round trips through conversion; fusing them into one crossing removes two. ## What varies between engines This whole cost exists only where there is a packed form to leave and return to. On an engine that holds ordinary language objects throughout, a same-language function receives the object the engine already had and pays no conversion whatsoever - its costs are the allocation rate and what the function itself does. Some engines have a packed form for one programming surface and objects for another, so the same logic can be free at one boundary and expensive at another in the same product. And the batching and shared-layout options above exist on some engines and not others. Establish which situation you are in before predicting the cost, and measure rather than assume. One adjacent effect is deliberately not this subject: once a step is an opaque function, the engine can no longer see inside it to rewrite the plan around it. That is a real consequence and a separate one. Here the concern is only the per-record cost of converting between the two representations.

  • How would you measure the boundary cost rather than infer it?
    Run the step twice: once with the real function, once with a function that returns its input unchanged. The second still pays every conversion and does no work, so the difference from the surrounding steps is the boundary, and the difference between the two runs is the logic. Comparing the step's records-in count with its duration tells you whether the cost is per record.
  • Why does placing the same function later in the pipeline often make it cheaper?
    Because the conversion is charged per record and per field materialised. Filtering and dropping unused fields upstream reduces both. Moving a per-record function to after a selective filter can cut its cost by the filter's selectivity without changing a line of its logic.
  • A result is kept for reuse in packed form and read by a per-record function each time. What is the trade?
    Keeping it packed makes it several times smaller, so more of it fits and less is displaced from other work. But every read through the function decodes it again, so the saving in footprint is paid for in repeated conversion. If the result is read many times through host-language code, keeping it as objects can be the cheaper of the two.

saying these in an interview costs you the question

  • Optimises the function's logic when the conversion is the cost.
  • Assumes a cross-language function runs inside the worker process.
  • Forgets that a second process's memory counts against the platform's ceiling.
  • Believes every engine pays a conversion for a user-supplied function.
  • Thinks per-record and per-batch crossings cost the same per record.
  • Places the crossing before the filters instead of after them.