adriangb commented on code in PR #24086:
URL: https://github.com/apache/datafusion/pull/24086#discussion_r4108774558
##########
datafusion/datasource-parquet/src/push_decoder.rs:
##########
@@ -721,6 +749,431 @@ impl PushDecoderStreamState {
}
}
+// `DF_FETCH_POLICY=streaming`: a batch-granular push decoder
+// (`FetchGranularity::Batch`) plus [`ReadAhead`], which fetches
+// [`ParquetPushDecoder::scan_plan`] ranges within a byte window.
+
+/// Peak bytes one streaming scan held (decoder buffers plus read-ahead in
+/// flight). Benchmarks read and reset it.
+pub static PEAK_STAGED_BYTES: std::sync::atomic::AtomicU64 =
+ std::sync::atomic::AtomicU64::new(0);
+
+/// How the stream schedules I/O.
+///
+/// - `Off`: fetch what the decoder asks for, a row group at a time.
+/// - `Streaming`: decode a batch at a time with up to `window` bytes of
+/// read-ahead.
+#[derive(Debug, Clone, Copy, PartialEq)]
+pub(crate) enum FetchPolicy {
+ Off,
+ Streaming { window: u64 },
+}
+
+impl FetchPolicy {
+ pub(crate) fn from_env() -> Self {
+ let window = std::env::var("DF_FETCH_BUDGET")
Review Comment:
Can we move this into a proper config?
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]