Skip to content

Add Arrow-native dataflow lane for ADBC sources and sinks - #814

Merged
flarco merged 1 commit into
v1.6.4from
flarco/fast-arrowlane
Sep 23, 2026
Merged

flarco merged 1 commit into
v1.6.4from
flarco/fast-arrowlane

Conversation

@flarco

@flarco flarco commented Sep 23, 2026

Copy link
Copy Markdown
Collaborator

Problem

Every database read builds a []any row per record and every write turns those rows back into the target format, even when both ends already speak Arrow. That extra hop costs CPU and allocations on large moves, and reading through the driver row API also drops type information the Arrow reader had (a Postgres numeric arrives as text, SQL Server datetimeoffset loses its offset, JSON columns lose their label).

Solution

Think of shipping sealed boxes instead of unpacking every item, carrying it across the room, and repacking it. When a source can hand over Arrow record batches and the target can take them, Sling now passes the batches straight through. Nothing about the row path changes when the lane does not apply: the run falls back to rows and logs why.

The lane

  • iop.RecordStream carries arrow.RecordBatch values with backpressure; iop.Datastream gained an Arrow mode (NewDatastreamArrow, ArrowOnly) that counts rows and bytes from records and never starts the row iterator.
  • The engine sits behind iop.ArrowLane (NewArrowLane). The open build ships a stub that declines, so an unlicensed build compiles and always takes the row path; the official build installs the real engine.
  • core/sling/task_run_arrow.go holds the eligibility gate: stage 1 checks the env switch, mode, source and target connections and the config; stage 2 runs once per stream on the real reader schema before the datastream starts, so the stream mode never changes after Start. SLING_ARROW_LANE accepts false (opt out), auto (default) and force (a decline fails the run).

Sources

  • ArrowDBConn reads records directly from ADBC; arrowLaneReader / laneExportStream / laneExportFlow are wired into Postgres, Snowflake, BigQuery and the generic BaseConn.BulkExportStream. A stage 2 decline returns ErrArrowLaneDeclined, so the caller reads with the native driver.
  • Driver type labels are used on the lane only (newAdbcLaneRead): Postgres numeric becomes a decimal column, jsonb a JSON column, uuid a UUID column. Driver extension fields are unwrapped before a file is written, so a Parquet file never carries an arrow.opaque type its storage cannot satisfy.
  • File sources (core/dbio/filesys/fs_arrow.go): NewArrowFileSet lists the files, reads every footer schema, and requires one shared schema before any record moves; a remote file is copied once to a temp file that the read then uses.
  • DuckDB connections with use_adbc now run every query, read and import through the ADBC handle (the instance file cannot be shared with the CLI session), and the ADBC path property is passed through to the driver.

Sinks

  • ADBC ingest and the staged-Parquet loaders take records as they arrive: stageFileFormat makes the staged COPY run as Parquet and stageDuckDbCompute skips the DuckDB merge for an Arrow dataflow; Redshift gets a copy_from_s3_parquet template.
  • File targets (filesys.writeDataflowRecords) stream records straight to the Parquet or Arrow IPC writer with the same part naming, partitioning, compression, readiness signals and byte counters as the row path. NewParquetArrowWriterFromSchema / WriteRecord and ArrowWriter.WriteRecord write whole records, with buffered row groups broken by counted bytes.
  • The row path is untouched: stageFileFormat returns the loader own default when the dataflow is not Arrow-only.

Supporting fixes

  • ArrowSchemaToColumns maps INT8/INT16 to SmallInt, TIME32/TIME64 to Time and DATE64 to Date; int16 builders were added to the Arrow and Parquet writers.
  • The Arrow and Parquet row readers no longer poll a done flag with a 10 ms sleep; they block on the row channel and the context.
  • Column stats gained a string update-key max (MaxStr) and are filled for Arrow streams from the record stream, so incremental state still advances.
  • ListFileNodes was extracted from ReadDataflow and shared with the Arrow file set.

Known limitations

  • The lane file-target write path hangs on this build, so suite entries 723-725 are wrapped in timeout 300 and fail fast instead of blocking the suite.
  • The cloud suite (tests/suite.cli.arrow.cloud.yaml, Databricks/Redshift) is manual; the Snowflake scale-0 read case stays there because the row path that follows the decline reads the table as 0 rows today.
  • A columns: spec always declines in this revision, so cases 743/755 assert the current decline text rather than a cannot cast reason.

Testing

  • New unit tests: iop record stream, datastream Arrow mode and metadata columns, Parquet/Arrow record writers, extension unwrapping; filesys Arrow file set and record writes; database ADBC lane read, stage loaders, ingest schema; sling gate table.
  • New CLI suite tests/suite.cli.arrow.yaml (included from tests/suite.cli.yaml), entries 701-791 covering the engage, fallback, quiet and force groups, with tests/pipelines/arrow/seed.yaml seeding 10,000 rows across the type matrix and verify.yaml comparing row-path and lane targets.
  • Run: just test-dbio-iop, just test-dbio-database, just test-dbio-filesys-local, just test-core, and cd cmd/sling && go build . && SLING_BIN=sling go test -v -run TestCLI -- --debug 700+.
  • justfile recipes now run their steps in subshells with set -e, so a failing step stops the recipe.

Wide tables spend most of a run building []any rows that are then
cast, merged and serialized right back. The lane carries Arrow record
batches instead: an ADBC reader feeds a RecordStream the sink pulls
from, so no row is materialized and no DuckDB merge runs.

- Gate eligibility before the read: config, connections, token and
  target columns decide, then the real reader schema is re-checked
  before the stream starts, so a stream never changes mode after
  Start. Every decline logs its reason and falls back to the row path.
- Targets: ADBC ingest, staged Parquet (Snowflake, Databricks,
  Redshift) and Parquet/Arrow files; Parquet/Arrow files also read as
  lane sources, with a shared-schema check over the file footers.
- Staged loaders and writers take records directly, keeping columns,
  metadata columns and string update keys so incremental state,
  stats and part naming match the row path.
- SLING_ARROW_LANE=false/force opt out or fail a decline; the open
  build keeps the stub engine and always takes the row path.
@flarco
flarco merged commit c1289d4 into v1.6.4 Sep 23, 2026
1 check passed
@flarco
flarco deleted the flarco/fast-arrowlane branch September 29, 2026 21:03
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant