Commit 7bde0dc
authored
Add Arrow IPC stream iterator for direct access to Arrow buffer (#279)
## Summary
This PR adds support for accessing raw Arrow IPC (Inter-Process
Communication) streams
directly, providing an alternative to the existing parsed Arrow record
interface. This
allows users to work with the Arrow binary format without the overhead
of parsing,
enabling use cases like streaming data to other systems or custom
processing pipelines.
## Changes
- **Refactored internal batch iterator** to expose IPC streams:
- Renamed `BatchIterator` → `IPCStreamIterator` (returns `io.Reader`)
- Updated implementations to return raw IPC streams instead of parsed
batches
- Created backward-compatible `BatchIterator` wrapper
- **Added public API** in `rows` package:
- New `GetArrowIPCStreams()` method on `Rows` interface
- New `ArrowIPCStreamIterator` interface with methods:
- `Next() (io.Reader, error)` - returns raw IPC stream
- `HasNext() bool` - `Close()` - `SchemaBytes() ([]byte, error)`
- **Updated all row scanner implementations** to support IPC streams
- **Added example** demonstrating IPC stream usage
## Benefits
- **Performance**: Avoid parsing overhead when forwarding Arrow data
- **Flexibility**: Direct access to Arrow binary format for custom
processing
- **Compatibility**: Easier integration with other Arrow-based systems
- **Memory efficiency**: Process streams without loading all records
into memory
## Testing
- All existing tests pass
- Backward compatibility maintained through wrapper pattern
- Example provided in `examples/ipcstreams/`
## Usage Example
```go
rows, err := db.QueryContext(ctx, "SELECT * FROM table")
ipcStreams, err := rows.(dbsqlrows.Rows).GetArrowIPCStreams(ctx)
defer ipcStreams.Close()
for ipcStreams.HasNext() {
reader, err := ipcStreams.Next()
// Process raw Arrow IPC stream
}
Signed-off-by: Jade Wang <jade.wang@databricks.com>1 parent 746c05d commit 7bde0dc
File tree
10 files changed
+455
-65
lines changed- examples/ipcstreams
- internal/rows
- arrowbased
- columnbased
- rowscanner
- rows
10 files changed
+455
-65
lines changed| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
| 1 | + | |
| 2 | + | |
| 3 | + | |
| 4 | + | |
| 5 | + | |
| 6 | + | |
| 7 | + | |
| 8 | + | |
| 9 | + | |
| 10 | + | |
| 11 | + | |
| 12 | + | |
| 13 | + | |
| 14 | + | |
| 15 | + | |
| 16 | + | |
| 17 | + | |
| 18 | + | |
| 19 | + | |
| 20 | + | |
| 21 | + | |
| 22 | + | |
| 23 | + | |
| 24 | + | |
| 25 | + | |
| 26 | + | |
| 27 | + | |
| 28 | + | |
| 29 | + | |
| 30 | + | |
| 31 | + | |
| 32 | + | |
| 33 | + | |
| 34 | + | |
| 35 | + | |
| 36 | + | |
| 37 | + | |
| 38 | + | |
| 39 | + | |
| 40 | + | |
| 41 | + | |
| 42 | + | |
| 43 | + | |
| 44 | + | |
| 45 | + | |
| 46 | + | |
| 47 | + | |
| 48 | + | |
| 49 | + | |
| 50 | + | |
| 51 | + | |
| 52 | + | |
| 53 | + | |
| 54 | + | |
| 55 | + | |
| 56 | + | |
| 57 | + | |
| 58 | + | |
| 59 | + | |
| 60 | + | |
| 61 | + | |
| 62 | + | |
| 63 | + | |
| 64 | + | |
| 65 | + | |
| 66 | + | |
| 67 | + | |
| 68 | + | |
| 69 | + | |
| 70 | + | |
| 71 | + | |
| 72 | + | |
| 73 | + | |
| 74 | + | |
| 75 | + | |
| 76 | + | |
| 77 | + | |
| 78 | + | |
| 79 | + | |
| 80 | + | |
| 81 | + | |
| 82 | + | |
| 83 | + | |
| 84 | + | |
| 85 | + | |
| 86 | + | |
| 87 | + | |
| 88 | + | |
| 89 | + | |
| 90 | + | |
| 91 | + | |
| 92 | + | |
| 93 | + | |
| 94 | + | |
| 95 | + | |
| 96 | + | |
| 97 | + | |
| 98 | + | |
| 99 | + | |
| 100 | + | |
| 101 | + | |
| 102 | + | |
| 103 | + | |
| 104 | + | |
| 105 | + | |
| 106 | + | |
| 107 | + | |
| 108 | + | |
| 109 | + | |
| 110 | + | |
| 111 | + | |
| 112 | + | |
| 113 | + | |
| 114 | + | |
| 115 | + | |
| 116 | + | |
| 117 | + | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
| 1 | + | |
| 2 | + | |
| 3 | + | |
| 4 | + | |
| 5 | + | |
| 6 | + | |
| 7 | + | |
| 8 | + | |
| 9 | + | |
| 10 | + | |
| 11 | + | |
| 12 | + | |
| 13 | + | |
| 14 | + | |
| 15 | + | |
| 16 | + | |
| 17 | + | |
| 18 | + | |
| 19 | + | |
| 20 | + | |
| 21 | + | |
| 22 | + | |
| 23 | + | |
| 24 | + | |
| 25 | + | |
| 26 | + | |
| 27 | + | |
| 28 | + | |
| 29 | + | |
| 30 | + | |
| 31 | + | |
| 32 | + | |
| 33 | + | |
| 34 | + | |
| 35 | + | |
| 36 | + | |
| 37 | + | |
| 38 | + | |
| 39 | + | |
| 40 | + | |
| 41 | + | |
| 42 | + | |
| 43 | + | |
| 44 | + | |
| 45 | + | |
| 46 | + | |
| 47 | + | |
| 48 | + | |
| 49 | + | |
| 50 | + | |
| 51 | + | |
| 52 | + | |
| 53 | + | |
| 54 | + | |
| 55 | + | |
| 56 | + | |
| 57 | + | |
| 58 | + | |
| 59 | + | |
| 60 | + | |
| 61 | + | |
| 62 | + | |
| 63 | + | |
| 64 | + | |
| 65 | + | |
| 66 | + | |
| 67 | + | |
| 68 | + | |
| 69 | + | |
| 70 | + | |
| 71 | + | |
| 72 | + | |
| 73 | + | |
| 74 | + | |
| 75 | + | |
| 76 | + | |
| 77 | + | |
| 78 | + | |
| 79 | + | |
| 80 | + | |
| 81 | + | |
| 82 | + | |
| 83 | + | |
| 84 | + | |
| 85 | + | |
| 86 | + | |
| 87 | + | |
| 88 | + | |
| 89 | + | |
| 90 | + | |
| 91 | + | |
| 92 | + | |
| 93 | + | |
| 94 | + | |
| 95 | + | |
| 96 | + | |
| 97 | + | |
| 98 | + | |
| 99 | + | |
| 100 | + | |
| 101 | + | |
| 102 | + | |
| 103 | + | |
| 104 | + | |
| 105 | + | |
| 106 | + | |
| 107 | + | |
| 108 | + | |
| 109 | + | |
| 110 | + | |
| 111 | + | |
| 112 | + | |
| 113 | + | |
| 114 | + | |
| 115 | + | |
| 116 | + | |
| 117 | + | |
| 118 | + | |
| 119 | + | |
| 120 | + | |
| 121 | + | |
| 122 | + | |
| 123 | + | |
| 124 | + | |
| 125 | + | |
| 126 | + | |
| 127 | + | |
| 128 | + | |
| 129 | + | |
| 130 | + | |
| 131 | + | |
| 132 | + | |
| 133 | + | |
| 134 | + | |
| 135 | + | |
| 136 | + | |
| 137 | + | |
| 138 | + | |
| 139 | + | |
| 140 | + | |
| 141 | + | |
| 142 | + | |
| 143 | + | |
| 144 | + | |
| 145 | + | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
329 | 329 | | |
330 | 330 | | |
331 | 331 | | |
| 332 | + | |
| 333 | + | |
| 334 | + | |
| 335 | + | |
| 336 | + | |
| 337 | + | |
| 338 | + | |
| 339 | + | |
| 340 | + | |
| 341 | + | |
| 342 | + | |
| 343 | + | |
| 344 | + | |
| 345 | + | |
332 | 346 | | |
333 | 347 | | |
334 | 348 | | |
| |||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
1041 | 1041 | | |
1042 | 1042 | | |
1043 | 1043 | | |
1044 | | - | |
1045 | | - | |
1046 | | - | |
1047 | | - | |
| 1044 | + | |
| 1045 | + | |
1048 | 1046 | | |
1049 | 1047 | | |
1050 | 1048 | | |
| |||
1674 | 1672 | | |
1675 | 1673 | | |
1676 | 1674 | | |
1677 | | - | |
| 1675 | + | |
1678 | 1676 | | |
1679 | 1677 | | |
1680 | 1678 | | |
1681 | 1679 | | |
1682 | 1680 | | |
1683 | | - | |
| 1681 | + | |
1684 | 1682 | | |
1685 | | - | |
| 1683 | + | |
1686 | 1684 | | |
1687 | 1685 | | |
1688 | 1686 | | |
1689 | 1687 | | |
1690 | 1688 | | |
1691 | 1689 | | |
1692 | | - | |
| 1690 | + | |
1693 | 1691 | | |
1694 | 1692 | | |
1695 | 1693 | | |
1696 | | - | |
| 1694 | + | |
1697 | 1695 | | |
1698 | 1696 | | |
1699 | 1697 | | |
| |||
0 commit comments