Add Arrow-native dataflow lane for ADBC sources and sinks - #814
Merged
Merged
Conversation
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.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Problem
Every database read builds a
[]anyrow 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 Postgresnumericarrives as text, SQL Serverdatetimeoffsetloses 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.RecordStreamcarriesarrow.RecordBatchvalues with backpressure;iop.Datastreamgained an Arrow mode (NewDatastreamArrow,ArrowOnly) that counts rows and bytes from records and never starts the row iterator.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.goholds 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 afterStart.SLING_ARROW_LANEacceptsfalse(opt out),auto(default) andforce(a decline fails the run).Sources
ArrowDBConnreads records directly from ADBC;arrowLaneReader/laneExportStream/laneExportFloware wired into Postgres, Snowflake, BigQuery and the genericBaseConn.BulkExportStream. A stage 2 decline returnsErrArrowLaneDeclined, so the caller reads with the native driver.newAdbcLaneRead): Postgresnumericbecomes a decimal column,jsonba JSON column,uuida UUID column. Driver extension fields are unwrapped before a file is written, so a Parquet file never carries anarrow.opaquetype its storage cannot satisfy.core/dbio/filesys/fs_arrow.go):NewArrowFileSetlists 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.use_adbcnow run every query, read and import through the ADBC handle (the instance file cannot be shared with the CLI session), and the ADBCpathproperty is passed through to the driver.Sinks
stageFileFormatmakes the staged COPY run as Parquet andstageDuckDbComputeskips the DuckDB merge for an Arrow dataflow; Redshift gets acopy_from_s3_parquettemplate.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/WriteRecordandArrowWriter.WriteRecordwrite whole records, with buffered row groups broken by counted bytes.stageFileFormatreturns the loader own default when the dataflow is not Arrow-only.Supporting fixes
ArrowSchemaToColumnsmaps INT8/INT16 to SmallInt, TIME32/TIME64 to Time and DATE64 to Date; int16 builders were added to the Arrow and Parquet writers.doneflag with a 10 ms sleep; they block on the row channel and the context.MaxStr) and are filled for Arrow streams from the record stream, so incremental state still advances.ListFileNodeswas extracted fromReadDataflowand shared with the Arrow file set.Known limitations
timeout 300and fail fast instead of blocking the 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.columns:spec always declines in this revision, so cases 743/755 assert the current decline text rather than acannot castreason.Testing
ioprecord stream, datastream Arrow mode and metadata columns, Parquet/Arrow record writers, extension unwrapping;filesysArrow file set and record writes;databaseADBC lane read, stage loaders, ingest schema;slinggate table.tests/suite.cli.arrow.yaml(included fromtests/suite.cli.yaml), entries 701-791 covering the engage, fallback, quiet and force groups, withtests/pipelines/arrow/seed.yamlseeding 10,000 rows across the type matrix andverify.yamlcomparing row-path and lane targets.just test-dbio-iop,just test-dbio-database,just test-dbio-filesys-local,just test-core, andcd cmd/sling && go build . && SLING_BIN=sling go test -v -run TestCLI -- --debug 700+.justfilerecipes now run their steps in subshells withset -e, so a failing step stops the recipe.