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

lxy-9602 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-cpp.git


The following commit(s) were added to refs/heads/main by this push:
     new ead92078 fix(aggregate): preserve listagg result ownership (#281)
ead92078 is described below

commit ead920782a85f375ed0356117a52914c1a231841
Author: 小明同学 <[email protected]>
AuthorDate: Fri Sep 4 20:00:53 2026 +0800

    fix(aggregate): preserve listagg result ownership (#281)
---
 .../compact/aggregate/field_listagg_agg.h          | 17 ++++----
 .../compact/aggregate/field_listagg_agg_test.cpp   | 46 ++++++++++++++--------
 .../operation/data_evolution_split_read_test.cpp   |  2 +-
 src/paimon/core/table/source/table_read_test.cpp   |  5 ++-
 src/paimon/core/utils/field_mapping_test.cpp       | 13 +++---
 test/inte/write_and_read_inte_test.cpp             | 44 +++++++++++++++++++++
 6 files changed, 91 insertions(+), 36 deletions(-)

diff --git a/src/paimon/core/mergetree/compact/aggregate/field_listagg_agg.h 
b/src/paimon/core/mergetree/compact/aggregate/field_listagg_agg.h
index a6dc04b7..ea1707ea 100644
--- a/src/paimon/core/mergetree/compact/aggregate/field_listagg_agg.h
+++ b/src/paimon/core/mergetree/compact/aggregate/field_listagg_agg.h
@@ -77,18 +77,16 @@ class FieldListaggAgg : public FieldAggregator {
             return input_field;
         }
 
+        std::string result;
         if (distinct_) {
-            result_ = AggDistinctImpl(acc_str, in_str);
+            result = AggDistinctImpl(acc_str, in_str);
         } else {
-            // Build into a local string to avoid aliasing when acc_str points 
into result_
-            std::string new_result;
-            new_result.reserve(acc_str.size() + delimiter_.size() + 
in_str.size());
-            new_result.append(acc_str);
-            new_result.append(delimiter_);
-            new_result.append(in_str);
-            result_ = std::move(new_result);
+            result.reserve(acc_str.size() + delimiter_.size() + in_str.size());
+            result.append(acc_str);
+            result.append(delimiter_);
+            result.append(in_str);
         }
-        return VariantType(std::string_view{result_});
+        return VariantType(BinaryString::FromString(result, pool_.get()));
     }
 
  private:
@@ -136,6 +134,5 @@ class FieldListaggAgg : public FieldAggregator {
 
     std::string delimiter_;
     bool distinct_;
-    std::string result_;
 };
 }  // namespace paimon
diff --git 
a/src/paimon/core/mergetree/compact/aggregate/field_listagg_agg_test.cpp 
b/src/paimon/core/mergetree/compact/aggregate/field_listagg_agg_test.cpp
index 1efac34e..01131847 100644
--- a/src/paimon/core/mergetree/compact/aggregate/field_listagg_agg_test.cpp
+++ b/src/paimon/core/mergetree/compact/aggregate/field_listagg_agg_test.cpp
@@ -47,13 +47,13 @@ class FieldListaggAggTest : public testing::Test {
 TEST_F(FieldListaggAggTest, TestSimple) {
     ASSERT_OK_AND_ASSIGN(auto agg, MakeAgg());
     auto ret = agg->Agg(std::string_view("hello"), std::string_view(" 
world")).value();
-    ASSERT_EQ(DataDefine::GetVariantValue<std::string_view>(ret), "hello, 
world");
+    ASSERT_EQ(DataDefine::GetStringView(ret), "hello, world");
 }
 
 TEST_F(FieldListaggAggTest, TestDelimiter) {
     ASSERT_OK_AND_ASSIGN(auto agg, MakeAgg("-"));
     auto ret = agg->Agg(std::string_view("user1"), 
std::string_view("user2")).value();
-    ASSERT_EQ(DataDefine::GetVariantValue<std::string_view>(ret), 
"user1-user2");
+    ASSERT_EQ(DataDefine::GetStringView(ret), "user1-user2");
 }
 
 TEST_F(FieldListaggAggTest, TestNull) {
@@ -62,12 +62,12 @@ TEST_F(FieldListaggAggTest, TestNull) {
     // input null -> return accumulator
     {
         auto ret = agg->Agg(std::string_view("hello"), NullType()).value();
-        ASSERT_EQ(DataDefine::GetVariantValue<std::string_view>(ret), "hello");
+        ASSERT_EQ(DataDefine::GetStringView(ret), "hello");
     }
     // accumulator null -> return input
     {
         auto ret = agg->Agg(NullType(), std::string_view("world")).value();
-        ASSERT_EQ(DataDefine::GetVariantValue<std::string_view>(ret), "world");
+        ASSERT_EQ(DataDefine::GetStringView(ret), "world");
     }
     // both null -> return null
     {
@@ -82,17 +82,17 @@ TEST_F(FieldListaggAggTest, TestEmptyString) {
     // empty input -> return accumulator
     {
         auto ret = agg->Agg(std::string_view("hello"), 
std::string_view("")).value();
-        ASSERT_EQ(DataDefine::GetVariantValue<std::string_view>(ret), "hello");
+        ASSERT_EQ(DataDefine::GetStringView(ret), "hello");
     }
     // empty accumulator -> return input
     {
         auto ret = agg->Agg(std::string_view(""), 
std::string_view("world")).value();
-        ASSERT_EQ(DataDefine::GetVariantValue<std::string_view>(ret), "world");
+        ASSERT_EQ(DataDefine::GetStringView(ret), "world");
     }
     // blank input -> return accumulator (which is empty)
     {
         auto ret = agg->Agg(std::string_view(""), 
std::string_view("")).value();
-        ASSERT_EQ(DataDefine::GetVariantValue<std::string_view>(ret), "");
+        ASSERT_EQ(DataDefine::GetStringView(ret), "");
     }
 }
 
@@ -112,12 +112,12 @@ TEST_F(FieldListaggAggTest, TestBlankStrings) {
                                                     u8" \t\u3000\u2000\n"};
     for (const std::string& blank : blank_strings) {
         auto ret = agg->Agg(std::string_view("user1"), 
std::string_view(blank)).value();
-        ASSERT_EQ(DataDefine::GetVariantValue<std::string_view>(ret), "user1");
+        ASSERT_EQ(DataDefine::GetStringView(ret), "user1");
     }
 
     // A blank accumulator must not add a leading delimiter.
     auto ret = agg->Agg(std::string_view(u8"\u3000\t"), 
std::string_view("user1")).value();
-    ASSERT_EQ(DataDefine::GetVariantValue<std::string_view>(ret), "user1");
+    ASSERT_EQ(DataDefine::GetStringView(ret), "user1");
 
     // A blank input must not turn a null accumulator into a non-null value.
     ret = agg->Agg(NullType(), std::string_view(u8" \t\u3000")).value();
@@ -129,9 +129,21 @@ TEST_F(FieldListaggAggTest, TestMultipleAccumulation) {
 
     // "a" + "," + "b" = "a,b", then "a,b" + "," + "c" = "a,b,c"
     auto ret = agg->Agg(std::string_view("a"), std::string_view("b")).value();
-    ASSERT_EQ(DataDefine::GetVariantValue<std::string_view>(ret), "a,b");
+    ASSERT_EQ(DataDefine::GetStringView(ret), "a,b");
     ret = agg->Agg(std::move(ret), std::string_view("c")).value();
-    ASSERT_EQ(DataDefine::GetVariantValue<std::string_view>(ret), "a,b,c");
+    ASSERT_EQ(DataDefine::GetStringView(ret), "a,b,c");
+}
+
+TEST_F(FieldListaggAggTest, TestResultOwnershipAcrossAggregations) {
+    ASSERT_OK_AND_ASSIGN(auto agg, MakeAgg());
+
+    ASSERT_OK_AND_ASSIGN(VariantType first,
+                         agg->Agg(std::string_view("alpha"), 
std::string_view("beta")));
+    ASSERT_OK_AND_ASSIGN(VariantType second,
+                         agg->Agg(std::string_view("one"), 
std::string_view("two")));
+
+    ASSERT_EQ(DataDefine::GetStringView(first), "alpha,beta");
+    ASSERT_EQ(DataDefine::GetStringView(second), "one,two");
 }
 
 TEST_F(FieldListaggAggTest, TestDistinct) {
@@ -139,7 +151,7 @@ TEST_F(FieldListaggAggTest, TestDistinct) {
 
     // "a;b" + "b;c" -> "a;b;c" (deduplicate "b")
     auto ret = agg->Agg(std::string_view("a;b"), 
std::string_view("b;c")).value();
-    ASSERT_EQ(DataDefine::GetVariantValue<std::string_view>(ret), "a;b;c");
+    ASSERT_EQ(DataDefine::GetStringView(ret), "a;b;c");
 }
 
 TEST_F(FieldListaggAggTest, TestDistinctIgnoresBlankTokens) {
@@ -148,7 +160,7 @@ TEST_F(FieldListaggAggTest, TestDistinctIgnoresBlankTokens) 
{
     auto ret =
         agg->Agg(std::string_view("user1"), std::string_view(u8" 
,user2,\t,\u3000,user1,\u2000"))
             .value();
-    ASSERT_EQ(DataDefine::GetVariantValue<std::string_view>(ret), 
"user1,user2");
+    ASSERT_EQ(DataDefine::GetStringView(ret), "user1,user2");
 }
 
 TEST_F(FieldListaggAggTest, TestDistinctNoDuplicates) {
@@ -156,7 +168,7 @@ TEST_F(FieldListaggAggTest, TestDistinctNoDuplicates) {
 
     // "a b" + "c d" -> "a b c d" (no dups to remove)
     auto ret = agg->Agg(std::string_view("a b"), std::string_view("c 
d")).value();
-    ASSERT_EQ(DataDefine::GetVariantValue<std::string_view>(ret), "a b c d");
+    ASSERT_EQ(DataDefine::GetStringView(ret), "a b c d");
 }
 
 TEST_F(FieldListaggAggTest, 
TestDistinctWithEmptyDelimiterFallsBackToWhitespace) {
@@ -164,7 +176,7 @@ TEST_F(FieldListaggAggTest, 
TestDistinctWithEmptyDelimiterFallsBackToWhitespace)
 
     // Empty delimiter falls back to whitespace, so the repeated "b" is 
removed.
     auto ret = agg->Agg(std::string_view("a b"), std::string_view("b 
c")).value();
-    ASSERT_EQ(DataDefine::GetVariantValue<std::string_view>(ret), "a b c");
+    ASSERT_EQ(DataDefine::GetStringView(ret), "a b c");
 }
 
 TEST_F(FieldListaggAggTest, TestDistinctEmptyInput) {
@@ -172,7 +184,7 @@ TEST_F(FieldListaggAggTest, TestDistinctEmptyInput) {
 
     // empty input -> return accumulator
     auto ret = agg->Agg(std::string_view("a;b"), std::string_view("")).value();
-    ASSERT_EQ(DataDefine::GetVariantValue<std::string_view>(ret), "a;b");
+    ASSERT_EQ(DataDefine::GetStringView(ret), "a;b");
 }
 
 TEST_F(FieldListaggAggTest, TestDistinctFalse) {
@@ -180,7 +192,7 @@ TEST_F(FieldListaggAggTest, TestDistinctFalse) {
 
     // "a;b" + "b;c" -> "a;b;b;c" (no dedup)
     auto ret = agg->Agg(std::string_view("a;b"), 
std::string_view("b;c")).value();
-    ASSERT_EQ(DataDefine::GetVariantValue<std::string_view>(ret), "a;b;b;c");
+    ASSERT_EQ(DataDefine::GetStringView(ret), "a;b;b;c");
 }
 
 TEST_F(FieldListaggAggTest, TestInvalidType) {
diff --git a/src/paimon/core/operation/data_evolution_split_read_test.cpp 
b/src/paimon/core/operation/data_evolution_split_read_test.cpp
index 03bc95b6..6b4067b0 100644
--- a/src/paimon/core/operation/data_evolution_split_read_test.cpp
+++ b/src/paimon/core/operation/data_evolution_split_read_test.cpp
@@ -111,7 +111,7 @@ TEST_F(DataEvolutionSplitReadTest, 
TestCreatePushDownPredicate) {
     auto f1_predicate =
         PredicateBuilder::Equal(/*field_index=*/1, /*field_name=*/"f1", 
FieldType::INT, Literal(2));
     auto row_id_predicate = PredicateBuilder::Equal(
-        /*field_index=*/2, SpecialFields::RowId().Name(), FieldType::BIGINT, 
Literal(3l));
+        /*field_index=*/2, SpecialFields::RowId().Name(), FieldType::BIGINT, 
Literal(int64_t{3}));
     ASSERT_OK_AND_ASSIGN(std::shared_ptr<Predicate> predicate,
                          PredicateBuilder::And({f0_predicate, f1_predicate, 
row_id_predicate}));
 
diff --git a/src/paimon/core/table/source/table_read_test.cpp 
b/src/paimon/core/table/source/table_read_test.cpp
index bc2c7f6f..f2a8644e 100644
--- a/src/paimon/core/table/source/table_read_test.cpp
+++ b/src/paimon/core/table/source/table_read_test.cpp
@@ -19,6 +19,7 @@
 
 #include "paimon/table/source/table_read.h"
 
+#include <cstdint>
 #include <map>
 #include <memory>
 #include <string>
@@ -52,7 +53,7 @@ TEST(TableReadTest, TestReadWithInvalidContext) {
     {
         // field type and literal type mismatch
         auto predicate = PredicateBuilder::Equal(/*field_index=*/3, 
/*field_name=*/"f3",
-                                                 FieldType::DOUBLE, 
Literal(15l));
+                                                 FieldType::DOUBLE, 
Literal(int64_t{15}));
         ReadContextBuilder context_builder(path);
         context_builder.SetPredicate(predicate);
         ASSERT_OK_AND_ASSIGN(auto read_context, context_builder.Finish());
@@ -63,7 +64,7 @@ TEST(TableReadTest, TestReadWithInvalidContext) {
     {
         // field type in predicate mismatch schema
         auto predicate = PredicateBuilder::Equal(/*field_index=*/3, 
/*field_name=*/"f3",
-                                                 FieldType::BIGINT, 
Literal(15l));
+                                                 FieldType::BIGINT, 
Literal(int64_t{15}));
         ReadContextBuilder context_builder(path);
         context_builder.SetPredicate(predicate);
         ASSERT_OK_AND_ASSIGN(auto read_context, context_builder.Finish());
diff --git a/src/paimon/core/utils/field_mapping_test.cpp 
b/src/paimon/core/utils/field_mapping_test.cpp
index 2de53d25..1ba674cd 100644
--- a/src/paimon/core/utils/field_mapping_test.cpp
+++ b/src/paimon/core/utils/field_mapping_test.cpp
@@ -20,6 +20,7 @@
 #include "paimon/core/utils/field_mapping.h"
 
 #include <algorithm>
+#include <cstdint>
 
 #include "arrow/type_fwd.h"
 #include "gtest/gtest.h"
@@ -444,7 +445,7 @@ TEST_F(FieldMappingTest, TestSchemaEvolutionWithPredicate) {
     auto greater_or_equal = PredicateBuilder::GreaterOrEqual(
         /*field_index=*/0, /*field_name=*/"key0", FieldType::INT, Literal(4));
     auto equal = PredicateBuilder::Equal(/*field_index=*/1, 
/*field_name=*/"key1",
-                                         FieldType::BIGINT, Literal(3l));
+                                         FieldType::BIGINT, 
Literal(int64_t{3}));
     auto less_or_equal = PredicateBuilder::LessOrEqual(/*field_index=*/2, 
/*field_name=*/"k",
                                                        FieldType::INT, 
Literal(10));
     // greater_than will not be pushed down, as with casting, only integer 
predicates can be pushed
@@ -455,7 +456,7 @@ TEST_F(FieldMappingTest, TestSchemaEvolutionWithPredicate) {
                                                 FieldType::INT, Literal(40));
     // in can be pushed down
     auto in = PredicateBuilder::In(/*field_index=*/5, /*field_name=*/"a", 
FieldType::BIGINT,
-                                   {Literal(100l)});
+                                   {Literal(int64_t{100})});
     auto not_in =
         PredicateBuilder::In(/*field_index=*/6, /*field_name=*/"e", 
FieldType::INT, {Literal(50)});
 
@@ -545,7 +546,7 @@ TEST_F(FieldMappingTest, TestSchemaEvolutionWithPredicate2) 
{
     auto greater_or_equal = PredicateBuilder::GreaterOrEqual(
         /*field_index=*/6, /*field_name=*/"key0", FieldType::INT, Literal(4));
     auto equal = PredicateBuilder::Equal(/*field_index=*/1, 
/*field_name=*/"key1",
-                                         FieldType::BIGINT, Literal(3l));
+                                         FieldType::BIGINT, 
Literal(int64_t{3}));
     auto less_or_equal = PredicateBuilder::LessOrEqual(/*field_index=*/4, 
/*field_name=*/"k",
                                                        FieldType::INT, 
Literal(10));
     // greater_than will not be pushed down, as with casting, only integer 
predicates can be pushed
@@ -556,7 +557,7 @@ TEST_F(FieldMappingTest, TestSchemaEvolutionWithPredicate2) 
{
                                                 FieldType::INT, Literal(40));
     // in will not be pushed down, as with casting, literal from BIGINT to INT 
is overflow
     auto in = PredicateBuilder::In(/*field_index=*/2, /*field_name=*/"a", 
FieldType::BIGINT,
-                                   {Literal(9223372036854775807l)});
+                                   {Literal(int64_t{9223372036854775807LL})});
     auto not_in =
         PredicateBuilder::In(/*field_index=*/3, /*field_name=*/"e", 
FieldType::INT, {Literal(50)});
 
@@ -577,7 +578,7 @@ TEST_F(FieldMappingTest, TestSchemaEvolutionWithPredicate2) 
{
     auto greater_or_equal_new = PredicateBuilder::GreaterOrEqual(
         /*field_index=*/0, /*field_name=*/"key0", FieldType::INT, Literal(4));
     auto equal_new = PredicateBuilder::Equal(/*field_index=*/1, 
/*field_name=*/"key1",
-                                             FieldType::BIGINT, Literal(3l));
+                                             FieldType::BIGINT, 
Literal(int64_t{3}));
     expected_part_info.partition_filter =
         PredicateBuilder::And({greater_or_equal_new, 
equal_new}).value_or(nullptr);
     CheckPartitionInfo(mapping->partition_info.value(), expected_part_info);
@@ -623,7 +624,7 @@ TEST_F(FieldMappingTest, 
TestCompoundPredicateWithoutPushDown) {
     auto equal = PredicateBuilder::Equal(/*field_index=*/0, 
/*field_name=*/"f0", FieldType::INT,
                                          Literal(55));
     auto greater_than = PredicateBuilder::GreaterThan(/*field_index=*/2, 
/*field_name=*/"f2",
-                                                      FieldType::BIGINT, 
Literal(30l));
+                                                      FieldType::BIGINT, 
Literal(int64_t{30}));
     auto less_than = PredicateBuilder::LessThan(/*field_index=*/3, 
/*field_name=*/"f3",
                                                 FieldType::INT, Literal(55));
     ASSERT_OK_AND_ASSIGN(auto or_predicate, 
PredicateBuilder::Or({greater_than, less_than}));
diff --git a/test/inte/write_and_read_inte_test.cpp 
b/test/inte/write_and_read_inte_test.cpp
index 43d616e8..c5d9b9d2 100644
--- a/test/inte/write_and_read_inte_test.cpp
+++ b/test/inte/write_and_read_inte_test.cpp
@@ -679,6 +679,50 @@ TEST_P(WriteAndReadInteTest, TestPKSimple) {
     ASSERT_TRUE(success);
 }
 
+TEST_P(WriteAndReadInteTest, TestPKListAggPreservesResultsAcrossKeys) {
+    arrow::FieldVector fields = {arrow::field("pk", arrow::utf8()),
+                                 arrow::field("value", arrow::utf8())};
+    auto [file_format, file_system] = GetParam();
+    std::map<std::string, std::string> options = {
+        {Options::FILE_FORMAT, file_format},
+        {Options::TARGET_FILE_SIZE, "1024"},
+        {Options::BUCKET, "1"},
+        {Options::FILE_SYSTEM, file_system},
+        {Options::MERGE_ENGINE, "aggregation"},
+        {"fields.value.aggregate-function", "listagg"},
+    };
+    if (file_system == "jindo") {
+        options = AddOptionsForJindo(options);
+    }
+    ASSERT_OK_AND_ASSIGN(
+        std::unique_ptr<TestHelper> helper,
+        TestHelper::Create(test_dir_, arrow::schema(fields), 
/*partition_keys=*/{},
+                           /*primary_keys=*/{"pk"}, options, 
/*is_streaming_mode=*/true));
+
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> first_batch,
+                         TestHelper::MakeRecordBatch(arrow::struct_(fields),
+                                                     R"([["first", "alpha"], 
["second", "one"]])",
+                                                     /*partition_map=*/{}, 
/*bucket=*/0, {}));
+    ASSERT_OK(helper->WriteAndCommit(std::move(first_batch), 
/*commit_identifier=*/0,
+                                     
/*expected_commit_messages=*/std::nullopt));
+    ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> second_batch,
+                         TestHelper::MakeRecordBatch(arrow::struct_(fields),
+                                                     R"([["first", "beta"], 
["second", "two"]])",
+                                                     /*partition_map=*/{}, 
/*bucket=*/0, {}));
+    ASSERT_OK(helper->WriteAndCommit(std::move(second_batch), 
/*commit_identifier=*/1,
+                                     
/*expected_commit_messages=*/std::nullopt));
+
+    arrow::FieldVector result_fields = fields;
+    result_fields.insert(result_fields.begin(), arrow::field("_VALUE_KIND", 
arrow::int8()));
+    ASSERT_OK_AND_ASSIGN(std::vector<std::shared_ptr<Split>> data_splits,
+                         helper->NewScan(StartupMode::LatestFull(), 
/*snapshot_id=*/std::nullopt));
+    ASSERT_OK_AND_ASSIGN(bool success,
+                         
helper->ReadAndCheckResult(arrow::struct_(result_fields), data_splits,
+                                                    R"([[0, "first", 
"alpha,beta"],
+                                       [0, "second", "one,two"]])"));
+    ASSERT_TRUE(success);
+}
+
 TEST_P(WriteAndReadInteTest, TestInputChangelogStreamRead) {
     arrow::FieldVector fields = {
         arrow::field("pk", arrow::utf8()),

Reply via email to