Est.

Foreign Function Interface in Columnar Analytics Engines

How columnar engines lose speed when calling foreign functions at row-oriented boundaries.

Staff Writer · · 10 min read
Cover illustration for “Foreign Function Interface in Columnar Analytics Engines”
Foreign Function and Runtime Extension Interfaces · October 7, 2026 · 10 min read · 2,255 words

A columnar analytics engine is fast because of a chain of properties that reinforce each other, and that chain snaps the moment a foreign function interface forces data across a row-oriented boundary. This piece traces that chain, prices out the cost of breaking it, and walks through how Apache Arrow, along with the FFI patterns built on top of it, try to hold the chain together.

What makes columnar engines fast, and what threatens it at the boundary

Storing data column by column instead of row by row looks like a minor layout choice with outsized consequences. When a table sits on disk organized by column, a query that only needs three columns out of fifty reads only those three. Nothing else gets pulled off disk. That alone cuts I/O by an order that scales with how wide the table is and how narrow the query's needs are.

Compression compounds that gain. Each technique works because the values sitting next to each other in memory are the same type and often similar in value, something a row can never offer.

Then there's the processing itself. Once data sits column-wise in memory, a single CPU instruction can act on thousands of values in one pass, through SIMD (single instruction, multiple data). Engines built around this layout process data in columnar batches, so CPU cycles go toward actual computation.

None of these three properties, reduced I/O, aggressive compression, and vectorized execution, work in isolation. Each one makes the others better. That compounding is the entire reason columnar engines outrun row stores on analytical workloads by such wide margins.

But what if the engine needs to run code someone else wrote? That's where the chain meets its hardest test. The instant a query calls out to a user-defined function sitting outside the engine's native execution path, something has to give. Columns get unpacked back into rows so the foreign code can read them one record at a time. Compression is gone, because the data has already been decoded into full-size values. SIMD is gone, because there's no longer a column to vectorize over. The I/O savings earned earlier in the pipeline don't disappear, but they stop mattering, because the bottleneck has moved downstream to a boundary crossing that erases everything the columnar layout bought.

What the serialization tax costs, mechanistically

Call this the serialization tax, and treat it as three separate penalties rather than one lump cost, because each penalty has its own mechanism and its own fix, or lack of one.

The first penalty is format conversion. That transposition is pure overhead: no new information gets created, no useful work gets done, the bytes just get rearranged into a shape the receiving code understands. This step alone cancels out the I/O advantage the engine fought to build, because the cost of rearranging the data can rival or exceed the cost of reading it off disk.

The second penalty is memory copying. The same values now live in memory twice, once in the engine's columnar buffers and once in the function's row-oriented copy. That duplication costs memory bandwidth and allocator time, and it scales with the size of the data crossing the boundary, so larger batches get hit harder, not easier.

The third penalty is the loss of vectorization. Once data has been flattened into rows, SIMD batch operations no longer apply. Each of these three penalties stacks on the ones before it, and the stacking happens at every single boundary crossing in a pipeline, not once per query.

These penalties cause slower queries in production, not just worse numbers in benchmarks. A user-defined function that adds one piece of business logic to an otherwise fast analytical query can leave that query slower than the equivalent row-store query would have been, because the columnar advantage gets consumed in its entirety at the boundary the UDF introduces. Any FFI approach that claims to solve this problem has to address all three penalties together. Fixing the copy while still losing vectorization doesn't either.

Diagram: The Three-Penalty Serialization Tax. Visualizes: Visualize the three stacking penalties that hit a columnar engine every time a foreign function interface crosses a row-oriented boundary.

How Apache Arrow's in-memory format makes zero-copy FFI possible

Apache Arrow was built to close that gap by giving every piece of a data pipeline, regardless of language or engine, a shared in-memory columnar format. Instead of each system inventing its own internal layout and forcing conversions at every handoff, systems that adopt Arrow describe their data the same way in memory, before it ever needs to move.

Arrow is a language-agnostic specification for representing structured, table-like data in memory, with a type system built specifically for the needs of analytical databases and data frame libraries. Because the format is standardized rather than proprietary, two systems that both speak Arrow can hand data to each other at little to no cost. The data just sits in a buffer that both sides already know how to read.

Inside a single process, this has a direct and practical payoff. When a UDF runtime and the query engine share the same Arrow buffers, the UDF receives its input columns exactly as the engine stored them. The three penalties from the previous section, conversion, copying, and lost vectorization, don't get minimized here. They get avoided, because the step that would have triggered them never happens.

Arrow's adoption across the analytics tooling world (query engines, data frame libraries, and language runtimes alike) has made it something close to a common tongue for columnar data. Systems that agree on Arrow can pass data to one another without paying serialization overhead at every hop. Zero-copy FFI doesn't happen because one system is clever. It happens because both sides agreed, in advance, on the exact shape data will take in memory. Arrow is that agreement, written down and implemented widely enough that the agreement actually holds across languages and engines.

Three FFI patterns production engines use, and their tradeoffs

Arrow being the shared agreement doesn't mean there's one way to build FFI on top of it. Production engines have settled on three distinct patterns, and picking one means deciding which part of the serialization tax gets fully eliminated and which risk gets accepted in exchange.

The first pattern runs UDFs in-process, directly on Arrow data. Native Arrow UDFs operate on columnar buffers without ever converting inputs into row-oriented objects, so the columnar layout holds end to end, no unnecessary copies get made, and vectorized processing stays intact through the entire execution path. Apache DataFusion, an open-source analytical query engine written in Rust and built on Arrow's columnar memory format, implements this about as completely as any engine does. Its built-in scalar, window, and aggregate functions use the exact same API that user-defined functions use, so a UDF isn't a bolted-on escape hatch. The entire engine becomes the security surface.

The second pattern runs UDFs inside WebAssembly, trading some of that eliminated tax back for sandboxing. Wasm has become the leading mechanism for running user-supplied code inside a database engine without spinning up a separate process, because it isolates memory without the cost of a process switch. The WAF paper (Huang et al., ICDE 2025) measures this overhead directly, and the results depend heavily on workload shape.

The third pattern moves the zero-copy boundary outward, from inside the engine to the connection layer between client and server. ADBC offers a vendor-neutral, Arrow-native API that gets rid of the row-conversion overhead baked into legacy ODBC and JDBC drivers. Where ODBC and JDBC convert a columnar result set into rows for transport, then force the receiving application to convert those rows back into columns, ADBC keeps the Arrow buffers intact the entire way across the driver boundary. That reframes the question. Instead of asking how to avoid row conversion inside one engine, ADBC asks how to avoid it across the whole data path, client to driver to server and back. It's a bigger problem to take on, but a more tractable one, because it only has to be solved once, at the transport layer.

Arroyo, a real-time SQL engine, handles the problem differently because Rust is statically typed and doesn't natively support loading arbitrary user code at runtime. Arroyo builds a dynamically linked, FFI-based plugin system to get dynamic UDF behavior out of a language that wasn't designed to offer it, a reminder that the three patterns above aren't the only shapes this problem can take, depending on what the host language will and won't let you do.

Where Arrow-native FFI still breaks down in practice

Arrow-native FFI removes the serialization tax along the path it was designed for, but a handful of real failure modes bring latency or risk back in through side doors.

Start with Wasm's warm-up gap. JIT compilation inside a Wasm runtime needs time to reach near-native throughput. A UDF that runs rarely, or only for a short burst, may never spend enough time executing to fully warm up.

Type system mismatches create a second failure mode. A UDF that emits a type Arrow doesn't represent cleanly can force the query planner out of its vectorized path entirely, pushing the rest of the plan into row-at-a-time processing and erasing columnar gains built up earlier in the query.

Security gaps make up a third failure mode, and this one sits apart from performance. Cloud data warehouses that support custom UDFs don't always pair that support with security mechanisms strong enough to manage the risk, or with the visibility an auditor would need to check what a UDF actually did. Wasm's memory isolation stops a UDF from corrupting engine memory, but it does nothing to stop timing attacks, I/O side channels, or data exfiltration carried out through results the query returns legitimately. The record of what a UDF touched, read, or sent elsewhere is often just missing.

A fourth failure mode occurs in the query optimizer. InterSystems IRIS 2025.2 takes this on directly, adding full support for mixing columnar and non-columnar indexes within a single query plan and updating the SQL optimizer's costing algorithm so columnar storage pays off more consistently in mixed transactional-analytical settings.

That leaves one question genuinely open. Does a JIT-warmed Wasm runtime paired with columnar batching outperform an out-of-process Arrow IPC transfer at realistic batch sizes? The WAF benchmarks test this directly, and the answer changes depending on workload shape. Anyone building on these patterns has to benchmark their own workload, because published numbers from one setup won't reliably predict another.

GPU-based query processing and the FFI problem at the hardware boundary

As analytical engines start pushing computation onto GPUs, a new version of the serialization tax appears at the boundary between CPU and GPU memory, and Arrow-native in-process FFI, built for CPU execution, has nothing to say about it.

It extends earlier tensor-computation-runtime systems that map relational operators, filters, joins, aggregations, onto vectorized tensor operations running on GPU hardware, and it does so by overlapping storage I/O, networking, and GPU computation.

The bottleneck PystachIO identifies is I/O overlap, not format conversion. Loading data from NVMe storage or over RDMA networks before computation starts creates a blocking execution model, one where the GPU and the available I/O bandwidth both sit underused while data trickles in, and where large storage-resident workloads can run into out-of-memory errors simply because the loading pattern wasn't designed with memory pressure in mind. PystachIO's answer is a form of pipeline parallelism: storage I/O, network transfer, and GPU computation overlap using chunked, asynchronous CUDA streams with deferred synchronization, so the GPU stays busy on one chunk of data while the next chunk is still arriving.

For FFI design, the implication is that the bottleneck moves. On CPU, Arrow makes format conversion unnecessary. At GPU scale, the harder problem becomes transfer scheduling and memory pressure management, concerns that need coordination at the engine level, not just agreement on a shared memory format. This remains an active research direction rather than a pattern already running in production systems, but it points toward where the next serialization tax is likely to appear as GPU infrastructure becomes more common in analytics stacks. The FFI problem isn't fixed in place. It moves as the hardware underneath it changes.

What agentic query patterns expose about FFI assumptions

Every FFI pattern covered so far was designed around a query pattern shaped by human behavior: a person types a query, waits, reads a result, and thinks before sending the next one. AI agents querying analytical engines don't behave that way. They issue queries continuously, often many in parallel, with no pause for a human to read a result before the next request fires.

That shift puts pressure on assumptions baked quietly into each FFI pattern discussed above. The instance pool sized for human query volume may simply be the wrong size for agentic volume.

The optimizer's cost-estimation blind spot around opaque UDF nodes raises a similar question under agentic load. The same mispricing, repeated across a continuous stream of agent-issued queries hitting the same UDF, compounds into a sustained cost the system has no mechanism to notice or correct.

Security gaps take on different weight too. The same gap, multiplied across a volume of agent-issued calls that no human is individually inspecting, raises the stakes on exactly the question the earlier section left unresolved: what a UDF actually did with the data it touched, and who would ever know.

None of this means the FFI patterns covered here are built wrong. It means they were built against a usage pattern that agentic systems don't follow, and the gap between those two patterns is where the next round of engineering work on columnar FFI is likely to concentrate.

Sources

  1. New in InterSystems IRIS 2025.2
  2. PystachIO: Efficient Distributed GPU Query Processing with PyTorch over Fast Networks & Fast Storage
  3. Apache Arrow
  4. Apache Arrow
  5. Leveraging Apache Arrow for Zero-copy, Zero-serialization Cluster Shared Memory
  6. Introducing the Apache Arrow C Data Interface

More in Foreign Function and Runtime Extension Interfaces