calvinchengx opened a new issue, #24044:
URL: https://github.com/apache/datafusion/issues/24044

   ## Describe the bug
   
   `RepartitionExec` coalesces on the producer side. `OutputChannel::coalesce`
   pushes each batch into a `LimitedBatchCoalescer` and forwards only batches 
that
   have reached `target_batch_size` (`SessionConfig::batch_size()`, default 
8192).
   Whatever is left over is flushed by `SharedCoalescer::finalize`, which runs 
when
   the last input sender finishes.
   
   An unbounded input never finishes, so the residual is never flushed. A stream
   that produces rows slowly, or in small batches, delivers nothing downstream
   until 8192 rows have accumulated — however long that takes. The plan is not
   deadlocked and no data is lost: it is buffered, invisibly, for an unbounded
   amount of time.
   
   The coalescer is created for every non-preserve-order repartition, without
   consulting the input's boundedness:
   
   
https://github.com/apache/datafusion/blob/main/datafusion/physical-plan/src/repartition/mod.rs#L528
   
   ```rust
   let shared_coalescer = (!preserve_order).then(|| {
       SharedCoalescer::new(
           input.schema(),
           context.session_config().batch_size(),
           num_input_partitions,
       )
   });
   ```
   
   `ExecutionPlan::boundedness()` reports `Boundedness::Unbounded` for such an
   input, but is not consulted here.
   
   ## To Reproduce
   
   Full program, depending only on `datafusion = "54.1.0"`, `futures`, and 
`tokio`.
   A `PartitionStream` emits one row every 100 ms and never ends; it is fed 
through
   `RepartitionExec` and the output is drained for six seconds.
   
   <details>
   <summary><code>src/main.rs</code></summary>
   
   ```rust
   //! `RepartitionExec` withholds all output from an unbounded input until
   //! `batch_size` rows accumulate.
   //!
   //! A source emits one row every 100 ms and never ends. It is fed through
   //! `RepartitionExec` with round-robin partitioning, and the output is 
drained
   //! for six seconds. About 60 rows are produced in that window.
   //!
   //! Run with `cargo run --release`.
   
   use std::sync::Arc;
   use std::time::{Duration, Instant};
   
   use datafusion::arrow::array::{Int64Array, RecordBatch};
   use datafusion::arrow::datatypes::{DataType, Field, Schema, SchemaRef};
   use datafusion::common::Result;
   use datafusion::execution::{SendableRecordBatchStream, TaskContext};
   use datafusion::physical_expr::Partitioning;
   use datafusion::physical_plan::repartition::RepartitionExec;
   use datafusion::physical_plan::stream::RecordBatchStreamAdapter;
   use datafusion::physical_plan::streaming::{PartitionStream, 
StreamingTableExec};
   use datafusion::physical_plan::{ExecutionPlan, execute_stream};
   use datafusion::physical_expr::LexOrdering;
   use datafusion::prelude::{SessionConfig, SessionContext};
   use futures::StreamExt;
   
   /// One row every 100 ms, forever.
   #[derive(Debug)]
   struct TickStream {
       schema: SchemaRef,
   }
   
   impl PartitionStream for TickStream {
       fn schema(&self) -> &SchemaRef {
           &self.schema
       }
   
       fn execute(&self, _ctx: Arc<TaskContext>) -> SendableRecordBatchStream {
           let schema = Arc::clone(&self.schema);
           let stream = futures::stream::unfold(0i64, move |i| {
               let schema = Arc::clone(&schema);
               async move {
                   tokio::time::sleep(Duration::from_millis(100)).await;
                   let batch = RecordBatch::try_new(
                       Arc::clone(&schema),
                       vec![Arc::new(Int64Array::from(vec![i]))],
                   );
                   Some((batch.map_err(Into::into), i + 1))
               }
           });
           Box::pin(RecordBatchStreamAdapter::new(
               Arc::clone(&self.schema),
               stream,
           ))
       }
   }
   
   /// Drains the plan for `secs` and returns how many rows arrived.
   async fn rows_delivered(batch_size: usize, secs: u64, repartition: bool) -> 
Result<usize> {
       let schema = Arc::new(Schema::new(vec![Field::new("v", DataType::Int64, 
false)]));
   
       let source = Arc::new(StreamingTableExec::try_new(
           Arc::clone(&schema),
           vec![Arc::new(TickStream {
               schema: Arc::clone(&schema),
           }) as Arc<dyn PartitionStream>],
           None,
           std::iter::empty::<LexOrdering>(),
           true, // infinite
           None,
       )?) as Arc<dyn ExecutionPlan>;
   
       let plan: Arc<dyn ExecutionPlan> = if repartition {
           Arc::new(RepartitionExec::try_new(
               source,
               Partitioning::RoundRobinBatch(4),
           )?)
       } else {
           source
       };
   
       let ctx = 
SessionContext::new_with_config(SessionConfig::new().with_batch_size(batch_size));
       let mut stream = execute_stream(plan, ctx.task_ctx())?;
   
       let deadline = Instant::now() + Duration::from_secs(secs);
       let mut rows = 0usize;
       while Instant::now() < deadline {
           match tokio::time::timeout(Duration::from_millis(250), 
stream.next()).await {
               Ok(Some(batch)) => rows += batch?.num_rows(),
               Ok(None) => break,
               Err(_) => {} // no batch this interval; keep waiting
           }
       }
       Ok(rows)
   }
   
   #[tokio::main]
   async fn main() -> Result<()> {
       let secs = 6;
       println!("source emits 1 row / 100 ms, so ~{} rows per run\n", secs * 
10);
   
       let baseline = rows_delivered(8192, secs, false).await?;
       println!("no repartition,   batch_size = 8192 -> {baseline} rows");
   
       let repart_default = rows_delivered(8192, secs, true).await?;
       println!("RoundRobinBatch(4), batch_size = 8192 -> {repart_default} 
rows");
   
       let repart_small = rows_delivered(1, secs, true).await?;
       println!("RoundRobinBatch(4), batch_size =    1 -> {repart_small} rows");
   
       Ok(())
   }
   ```
   
   </details>
   
   Output of `cargo run --release`:
   
   ```
   source emits 1 row / 100 ms, so ~60 rows per run
   
   no repartition,   batch_size = 8192 -> 59 rows
   RoundRobinBatch(4), batch_size = 8192 -> 0 rows
   RoundRobinBatch(4), batch_size =    1 -> 59 rows
   ```
   
   The middle line is the bug: the same source, the same six seconds, one
   `RepartitionExec` in between, and nothing arrives. The first line shows the
   source itself is fine, and the third isolates the coalescer — the 
repartition is
   still there, only `batch_size` changed.
   
   Reproduced on 54.1.0; the construction on `main` is unchanged.
   
   ## Expected behavior
   
   Rows reach the consumer at a bounded latency, as they do when the same plan 
runs
   without a repartition, or with a small `batch_size`.
   
   For an unbounded input, waiting for a full batch is not a bounded wait. A few
   directions, in the order I would guess you prefer them — though this is your
   design call, which is why this is an issue and not a pull request:
   
   1. Skip the producer-side coalescer when the input is 
`Boundedness::Unbounded`.
   2. Keep it, but flush on a time bound as well as a size bound.
   3. Expose it as a config so an unbounded pipeline can turn it off.
   
   ## Additional context
   
   Some notes to save duplicated work:
   
   - **This is not a regression from #18782.** The producer-side coalescer 
arrived
     with that issue, which also removed the `coalesce_batches` optimizer rule 
that
     previously inserted a `CoalesceBatchesExec` after `RepartitionExec`. I 
checked
     the older rule (52.0.0) and it did not consider boundedness either, so the
     blind spot predates the integration — #18782 moved it inside the operator
     rather than introducing it.
   - **`FilterExec` looks like it has the same shape** — it holds a
     `LimitedBatchCoalescer` and I did not see a boundedness check near it. I 
have
     **not** verified that one behaves the same way, so treat it as a lead 
rather
     than a second report.
   - Found while debugging a streaming engine built on DataFusion, where a
     round-robin repartition sat between an unbounded source and a sink and the
     query produced no output at all. Worked around downstream by not 
repartitioning
     streaming plans; the underlying behaviour is unchanged, which is why I am
     reporting it here.
   


-- 
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]

Reply via email to