feat(table): expose `registerStreamingTable` for push-mode batch ingest
Chưa có ai nhận issue này.
Đánh giá
- Độ khó
- 5/5
- Thời gian dự kiến
- Hơn một tuần
- Mức phù hợp với người mới
- 38/100
Hướng nghiên cứu
Bắt đầu với đường dẫn SessionContext.registerTable/TableProvider hiện có từ PR #65 và điểm vào của phần triển khai trước đó rust/src/api.rs:572, register_partition_stream. Theo dõi việc tích hợp Data.exportVectorSchemaRoot và StreamingTable/PartitionStream được mô tả trong issue; được xem là hoàn tất khi các producer có thể ghi, đóng hoặc thất bại đồng thời với các thao tác đọc của truy vấn, với backpressure, hủy, lan truyền lỗi và ngữ nghĩa single-scan được ghi lại và thực thi.
Do mô hình lập chỉ mục viết ra từ nội dung của issue.
Mô tả
Is your feature request related to a problem or challenge?
PR #65 shipped a Java-implemented TableProvider and SessionContext.registerTable(String, TableProvider). That covered the pull shape: DataFusion calls scan(BufferAllocator) and reads the returned ArrowReader.
But it does not cover the push shape that event-driven batch sources need:
- A coordinator that reduces over shard responses arriving incrementally -- the producer can't materialise an
ArrowReaderbecause the next batch hasn't arrived yet. - A Flight stream feeding into a query -- same problem; the producer is event-driven.
- Any in-process producer that emits batches as side-effects of other work and doesn't know in advance how many will arrive.
To bridge these into PR #65 today, callers have to write a BlockingArrowReader adapter that buffers pushed batches and serves them through the pull interface. That's a serialisation point: the producer blocks waiting for loadNextBatch() to be called, or DataFusion blocks waiting for the next batch -- the two ends can never run truly concurrently. The adapter also has to invent its own backpressure semantics, error propagation, end-of-stream signalling, and thread-safety story.
DataFusion itself solves this on the Rust side with StreamingTable + PartitionStream plus an mpsc channel: producer pushes Result<RecordBatch> into the sender, the consumer (DataFusion's StreamingTableExec) polls the receiver as part of normal query execution. The two ends decouple via the channel buffer, with the runtime providing backpressure and cancellation propagation.
Describe the solution you'd like
One new method on SessionContext returning a TableSink:
TableSink sink = ctx.registerStreamingTable("shard_results", schema, capacity);
// Producer thread (any thread, including outside any Tokio runtime):
try {
while (hasMoreInput()) {
sink.write(batch); // backpressures when channel is full
}
sink.close(); // EOF: queries see end-of-stream cleanly
} catch (Throwable t) {
sink.fail(t); // signal error: queries see RuntimeException
}
public final class TableSink implements AutoCloseable {
void write(VectorSchemaRoot batch); // exports via Data.exportVectorSchemaRoot
void close(); // EOF
void fail(Throwable cause); // error propagated to readers
}
After registration the table can be referenced like any other registered table:
DataFrame df = ctx.sql("SELECT count(*) FROM shard_results");
ArrowReader r = df.executeStream(allocator);
// Producer thread continues writing as r.loadNextBatch() drains.
Single-scan semantics. The registered table can only be queried once. After that scan completes (or is cancelled), the sink is no longer usable and the table cannot be re-scanned. This is the natural semantic for an event-driven producer -- the data is consumed as it arrives -- and matches what every downstream Substrait/Calcite plan that uses streaming tables already assumes. Documented loudly on registerStreamingTable's Javadoc; trying to re-execute against the same registration throws.
This is intentional. Re-scannable streaming would require buffering every batch internally, which defeats the streaming use case. Callers who need to re-scan the same data should use the existing registerTable / SimpleTableProvider pull shape (PR #65) instead.
Describe alternatives you've considered
BlockingArrowReaderadapter on top of PR #65'sregisterTable. What every caller currently has to hand-roll. It works but pushes the channel + backpressure + EOF + error story onto every embedder. Bridging via the upstream-canonicalStreamingTableshape is strictly less code and gets cancellation propagation for free.- Backpressure-free
try_write. A non-blocking variant that returnsfalsewhen the channel is full. Easy to add later as a follow-up if anyone wants it; not in scope here. Defaultwriteblocks, which is the contract every Java I/O caller expects. - Reuse PR #65's
TableProviderinterface and wrap mpsc internally. Considered. The problem: PR #65'sscan(BufferAllocator) -> ArrowReaderreturns synchronously, so a mpsc-backed implementation has to block onloadNextBatch()waiting for the producer -- exactly the serialisation point we're trying to avoid. Going direct toStreamingTable+PartitionStreamis the right layer.
Additional context
- The OpenSearch backend's
rust/src/api.rs:572register_partition_streamis the prior-art template; it does almost exactly this. The Java side there uses a hand-rolled FFM bridge (sender_send) that can be replaced with this surface as soon as it lands.
- Ngôn ngữ chính
- Java
- Star
- 32
- Fork
- 12
- Chỉ số merge pull request
- Không có pull request nào được merge trong 30 ngày
Hướng dẫn đóng góp
Bắt đầu từ đâu
- Đọc hết issue, rồi đọc hướng dẫn đóng góp của dự án.
- Bình luận trên issue rằng bạn sẽ nhận — tránh hai người làm cùng một việc.
- Fork repository và làm thay đổi trên một nhánh.
- Mở pull request có tham chiếu số hiệu của issue.
Issue khác của apache/datafusion-java
-
Độ khó 5/5 Hơn một tuần Mức phù hợp với người mới 35/100
apache/datafusion-java#116 ·
-
Độ khó 5/5 Hơn một tuần Mức phù hợp với người mới 25/100
apache/datafusion-java#112 ·
-
enhancement
Độ khó 5/5 Hơn một tuần Mức phù hợp với người mới 42/100
apache/datafusion-java#96 ·
-
Create first release Đang mởenhancement
Độ khó 4/5 3-5 ngày Mức phù hợp với người mới 35/100
apache/datafusion-java#86 · 3 bình luận ·
-
enhancement
Độ khó 5/5 Hơn một tuần Mức phù hợp với người mới 45/100
apache/datafusion-java#68 ·
Tất cả issue của apache/datafusion-java
Issue tương tự
-
bug untriaged
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 84/100
opensearch-project/ml-commons#5094 ·
-
bug
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 85/100
-
emitter:client:csharp feature
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 72/100
-
affects/8.10 affects/8.9 component/clients kind/bug likelihood/mid severity/mid
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 78/100
-
Two open-case totals on one screen: the Programs tile says 15,858 and the nav badge says 15,868 Đang mởbug frontend maui-pilot
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 72/100