This is an automated email from the ASF dual-hosted git repository.

zanmato1984 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow.git


The following commit(s) were added to refs/heads/main by this push:
     new 43eca861e59 GH-46454: [C++][Dataset][Acero] Preserve order when 
writting with TeeNode (#46455)
43eca861e59 is described below

commit 43eca861e59a85667c7cce05816b899da9ac2bef
Author: gitmodimo <[email protected]>
AuthorDate: Fri Aug 28 20:22:58 2026 +0200

    GH-46454: [C++][Dataset][Acero] Preserve order when writting with TeeNode 
(#46455)
    
    ### Rationale for this change
    TeeNode needs to sequence batches when  implicit order within processed 
dataset.
    
    ### What changes are included in this PR?
    Conditionally sequence batches when preserve_order=true
    
    ### Are these changes tested?
    I tested it in my use case. No CI tests AFAIK.
    
    ### Are there any user-facing changes?
    
    Dataset will now be ordered as expected.
    
    * GitHub Issue: #46454
    
    Authored-by: RafaƂ Hibner <[email protected]>
    Signed-off-by: Rossi Sun <[email protected]>
---
 cpp/src/arrow/dataset/file_base.cc |  33 +++++++++-
 cpp/src/arrow/dataset/file_test.cc | 130 ++++++++++++++++++++++++++++++++-----
 2 files changed, 145 insertions(+), 18 deletions(-)

diff --git a/cpp/src/arrow/dataset/file_base.cc 
b/cpp/src/arrow/dataset/file_base.cc
index ccc79dfa9bf..8ad2254aa64 100644
--- a/cpp/src/arrow/dataset/file_base.cc
+++ b/cpp/src/arrow/dataset/file_base.cc
@@ -17,6 +17,7 @@
 
 #include "arrow/dataset/file_base.h"
 
+#include "arrow/acero/accumulation_queue.h"
 #include "arrow/acero/exec_plan.h"
 
 #include <algorithm>
@@ -559,13 +560,18 @@ Result<acero::ExecNode*> MakeWriteNode(acero::ExecPlan* 
plan,
   return node;
 }
 
-class TeeNode : public acero::MapNode {
+class TeeNode : public acero::MapNode,
+                public arrow::acero::util::SerialSequencingQueue::Processor {
  public:
   TeeNode(acero::ExecPlan* plan, std::vector<acero::ExecNode*> inputs,
           std::shared_ptr<Schema> output_schema,
           FileSystemDatasetWriteOptions write_options)
       : MapNode(plan, std::move(inputs), std::move(output_schema)),
-        write_options_(std::move(write_options)) {}
+        write_options_(std::move(write_options)) {
+    if (write_options.preserve_order) {
+      sequencer_ = acero::util::SerialSequencingQueue::Make(this);
+    }
+  }
 
   Status StartProducing() override {
     ARROW_ASSIGN_OR_RAISE(
@@ -592,6 +598,28 @@ class TeeNode : public acero::MapNode {
 
   const char* kind_name() const override { return "TeeNode"; }
 
+  Status Validate() const override {
+    ARROW_RETURN_NOT_OK(acero::MapNode::Validate());
+    if (inputs_[0]->ordering().is_unordered() && sequencer_) {
+      return Status::Invalid("Tee node '", label(),
+                             "' is configured to sequence output but there is 
no "
+                             "meaningful ordering in the input");
+    }
+    return Status::OK();
+  }
+
+  Status InputReceived(ExecNode* input, ExecBatch batch) override {
+    DCHECK_EQ(input, inputs_[0]);
+    if (sequencer_) {
+      return sequencer_->InsertBatch(std::move(batch));
+    }
+    return Process(std::move(batch));
+  }
+
+  Status Process(ExecBatch batch) override {
+    return acero::MapNode::InputReceived(inputs_[0], batch);
+  }
+
   void Finish() override { dataset_writer_->Finish(); }
 
   Result<compute::ExecBatch> ProcessBatch(compute::ExecBatch batch) override {
@@ -625,6 +653,7 @@ class TeeNode : public acero::MapNode {
   std::unique_ptr<internal::DatasetWriter> dataset_writer_;
   FileSystemDatasetWriteOptions write_options_;
   std::atomic<int32_t> backpressure_counter_ = 0;
+  std::unique_ptr<acero::util::SerialSequencingQueue> sequencer_{nullptr};
 };
 
 }  // namespace
diff --git a/cpp/src/arrow/dataset/file_test.cc 
b/cpp/src/arrow/dataset/file_test.cc
index 2e2561203be..bd87a58eb41 100644
--- a/cpp/src/arrow/dataset/file_test.cc
+++ b/cpp/src/arrow/dataset/file_test.cc
@@ -32,6 +32,7 @@
 #include <arrow/record_batch.h>
 #include <arrow/util/async_generator.h>
 #include "arrow/acero/exec_plan.h"
+#include "arrow/acero/test_nodes.h"
 #include "arrow/acero/test_util_internal.h"
 #include "arrow/array/array_primitive.h"
 #include "arrow/compute/test_util_internal.h"
@@ -44,6 +45,7 @@
 #include "arrow/filesystem/test_util.h"
 #include "arrow/status.h"
 #include "arrow/testing/future_util.h"
+#include "arrow/testing/generator.h"
 #include "arrow/testing/gtest_util.h"
 #include "arrow/util/io_util.h"
 
@@ -444,6 +446,65 @@ class MockDataset : public Dataset {
   };
 };
 
+constexpr random::SeedType kJitterSeed = 42;
+constexpr int kMaxJitterModifier = 4;
+constexpr int64_t kOrderingRowsPerBatch = 1;
+constexpr int kOrderingNumBatches = 256;
+
+Result<bool> HasOutOfOrderRows(const Table& table) {
+  TableBatchReader reader(table);
+  std::shared_ptr<RecordBatch> batch;
+  ARROW_RETURN_NOT_OK(reader.ReadNext(&batch));
+  int32_t prev = 0;
+  bool has_prev = false;
+  while (batch != nullptr) {
+    const auto* values = batch->column(0)->data()->GetValues<int32_t>(1);
+    for (int row = 0; row < batch->num_rows(); ++row) {
+      int32_t value = values[row];
+      if (has_prev && value <= prev) {
+        return true;
+      }
+      prev = value;
+      has_prev = true;
+    }
+    ARROW_RETURN_NOT_OK(reader.ReadNext(&batch));
+  }
+  return false;
+}
+
+TEST_F(TestFileSystemDataset, RejectPreserveOrderWithUnorderedInput) {
+  dataset::internal::Initialize();
+
+  auto format = std::make_shared<IpcFileFormat>();
+  FileSystemDatasetWriteOptions write_options;
+  write_options.file_write_options = format->DefaultWriteOptions();
+  write_options.filesystem = 
std::make_shared<fs::internal::MockFileSystem>(fs::kNoTime);
+  write_options.base_dir = "root";
+  write_options.partitioning = std::make_shared<HivePartitioning>(schema({}));
+  write_options.basename_template = "{i}.feather";
+  write_options.preserve_order = true;
+
+  auto source_data = acero::MakeBasicBatches();
+  for (const char* factory_name : {"write", "tee"}) {
+    SCOPED_TRACE(factory_name);
+    ASSERT_OK_AND_ASSIGN(auto plan, acero::ExecPlan::Make());
+    AsyncGenerator<std::optional<cp::ExecBatch>> sink_gen;
+    std::vector<acero::Declaration> declarations = {
+        {"source",
+         acero::SourceNodeOptions{source_data.schema, source_data.gen(false, 
false)}},
+        {factory_name, WriteNodeOptions{write_options}},
+    };
+    if (std::string(factory_name) == "tee") {
+      declarations.emplace_back("sink", acero::SinkNodeOptions{&sink_gen});
+    }
+    ASSERT_OK(
+        
acero::Declaration::Sequence(std::move(declarations)).AddToPlan(plan.get()));
+    ASSERT_THAT(plan->Validate(),
+                Raises(StatusCode::Invalid,
+                       ::testing::HasSubstr("no meaningful ordering in the 
input")));
+  }
+}
+
 TEST_F(TestFileSystemDataset, MultiThreadedWritePersistsOrder) {
   // Test for GH-26818
   //
@@ -458,6 +519,7 @@ TEST_F(TestFileSystemDataset, 
MultiThreadedWritePersistsOrder) {
   //
   // If this test starts to reliably fail with preserve_order == false, the 
test setup
   // has to be revised to again reliably produce out-of-order sequences.
+
   auto format = std::make_shared<IpcFileFormat>();
   FileSystemDatasetWriteOptions write_options;
   write_options.file_write_options = format->DefaultWriteOptions();
@@ -500,26 +562,62 @@ TEST_F(TestFileSystemDataset, 
MultiThreadedWritePersistsOrder) {
     ASSERT_OK(scanner_builder->UseThreads(false));
     ASSERT_OK_AND_ASSIGN(scanner, scanner_builder->Finish());
     ASSERT_OK_AND_ASSIGN(auto actual, scanner->ToTable());
-    TableBatchReader reader(*actual);
-    std::shared_ptr<RecordBatch> batch;
-    ASSERT_OK(reader.ReadNext(&batch));
-    int32_t prev = -1;
-    auto out_of_order = false;
-    while (batch != nullptr) {
-      const auto* values = batch->column(0)->data()->GetValues<int32_t>(1);
-      for (int row = 0; row < batch->num_rows(); ++row) {
-        int32_t value = values[row];
-        if (value <= prev) {
-          out_of_order = true;
-        }
-        prev = value;
-      }
-      ASSERT_OK(reader.ReadNext(&batch));
-    }
+    ASSERT_OK_AND_ASSIGN(auto out_of_order, HasOutOfOrderRows(*actual));
     ASSERT_EQ(!out_of_order, preserve_order);
   }
 }
 
+TEST_F(TestFileSystemDataset, MultiThreadedTeeWritePersistsOrder) {
+  dataset::internal::Initialize();
+  acero::RegisterTestNodes();
+
+  auto format = std::make_shared<IpcFileFormat>();
+  auto fs = std::make_shared<fs::internal::MockFileSystem>(fs::kNoTime);
+  FileSystemDatasetWriteOptions write_options;
+  write_options.file_write_options = format->DefaultWriteOptions();
+  write_options.filesystem = fs;
+  write_options.partitioning = std::make_shared<HivePartitioning>(schema({}));
+  write_options.basename_template = "{i}.feather";
+
+  auto unordered_write_options = write_options;
+  unordered_write_options.base_dir = "unordered";
+  unordered_write_options.preserve_order = false;
+  auto ordered_write_options = write_options;
+  ordered_write_options.base_dir = "ordered";
+  ordered_write_options.preserve_order = true;
+
+  auto input = gen::Gen({gen::Step<int32_t>()})
+                   ->FailOnError()
+                   ->Table(kOrderingRowsPerBatch, kOrderingNumBatches);
+
+  // The first TeeNode records the jittered, out-of-order stream without 
changing it.
+  // The second TeeNode must use the batch indices to restore order.
+  ASSERT_OK(acero::DeclarationToStatus(acero::Declaration::Sequence(
+      {{"table_source", acero::TableSourceNodeOptions{input}},
+       {"jitter", acero::JitterNodeOptions{kJitterSeed, kMaxJitterModifier}},
+       {"tee", WriteNodeOptions{unordered_write_options}, "unordered_tee"},
+       {"tee", WriteNodeOptions{ordered_write_options}, "ordered_tee"}})));
+
+  auto read_written_table =
+      [&](const std::string& path) -> Result<std::shared_ptr<Table>> {
+    ARROW_ASSIGN_OR_RAISE(auto dataset_factory,
+                          FileSystemDatasetFactory::Make(fs, {path}, format, 
{}));
+    ARROW_ASSIGN_OR_RAISE(auto written_dataset, 
dataset_factory->Finish(FinishOptions{}));
+    ARROW_ASSIGN_OR_RAISE(auto written_scanner_builder, 
written_dataset->NewScan());
+    ARROW_RETURN_NOT_OK(written_scanner_builder->UseThreads(false));
+    ARROW_ASSIGN_OR_RAISE(auto written_scanner, 
written_scanner_builder->Finish());
+    return written_scanner->ToTable();
+  };
+
+  ASSERT_OK_AND_ASSIGN(auto unordered_table, 
read_written_table("unordered/0.feather"));
+  ASSERT_OK_AND_ASSIGN(auto unordered_out_of_order, 
HasOutOfOrderRows(*unordered_table));
+  ASSERT_TRUE(unordered_out_of_order);
+
+  ASSERT_OK_AND_ASSIGN(auto ordered_table, 
read_written_table("ordered/0.feather"));
+  ASSERT_OK_AND_ASSIGN(auto ordered_out_of_order, 
HasOutOfOrderRows(*ordered_table));
+  ASSERT_FALSE(ordered_out_of_order);
+}
+
 class FileSystemWriteTest : public testing::TestWithParam<std::tuple<bool, 
bool>> {
   using PlanFactory = std::function<std::vector<acero::Declaration>(
       const FileSystemDatasetWriteOptions&,

Reply via email to