Skip to content

Commit fe90ca1

Browse files
committed
fix review 4
1 parent dba2244 commit fe90ca1

8 files changed

Lines changed: 101 additions & 30 deletions

File tree

cmake_modules/DefineOptions.cmake

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -92,7 +92,7 @@ if("${CMAKE_SOURCE_DIR}" STREQUAL "${CMAKE_CURRENT_SOURCE_DIR}")
9292
define_option_string(PAIMON_CXXFLAGS "Compiler flags to append when compiling Paimon"
9393
"")
9494

95-
define_option(PAIMON_BUILD_STATIC "Build static libraries" OFF)
95+
define_option(PAIMON_BUILD_STATIC "Build static libraries" ON)
9696

9797
define_option(PAIMON_BUILD_SHARED "Build shared libraries" ON)
9898
#----------------------------------------------------------------------

include/paimon/defs.h

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -130,7 +130,7 @@ struct PAIMON_EXPORT Options {
130130
static const char MANIFEST_TARGET_FILE_SIZE[];
131131

132132
/// "manifest.format" - Specify the message format of manifest files.
133-
/// Default value is orc.
133+
/// Default value is avro.
134134
static const char MANIFEST_FORMAT[];
135135

136136
/// "manifest.compression" - File compression for manifest, default value is zstd.

src/paimon/common/utils/arrow/arrow_utils.h

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -39,14 +39,17 @@ class ArrowUtils {
3939
return std::make_shared<arrow::Schema>(struct_type->fields());
4040
}
4141

42-
static std::vector<int32_t> CreateProjection(
42+
static Result<std::vector<int32_t>> CreateProjection(
4343
const std::shared_ptr<::arrow::Schema>& file_schema,
4444
const arrow::FieldVector& read_fields) {
4545
std::vector<int32_t> target_to_src_mapping;
4646
target_to_src_mapping.reserve(read_fields.size());
4747
for (const auto& field : read_fields) {
4848
auto src_field_idx = file_schema->GetFieldIndex(field->name());
49-
assert(src_field_idx >= 0);
49+
if (src_field_idx < 0) {
50+
return Status::Invalid(
51+
fmt::format("Field '{}' not found or duplicate in file schema", field->name()));
52+
}
5053
target_to_src_mapping.push_back(src_field_idx);
5154
}
5255
return target_to_src_mapping;

src/paimon/common/utils/arrow/arrow_utils_test.cpp

Lines changed: 82 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -18,30 +18,93 @@
1818

1919
#include "arrow/api.h"
2020
#include "gtest/gtest.h"
21-
#include "paimon/common/types/data_field.h"
21+
#include "paimon/testing/utils/testharness.h"
2222

2323
namespace paimon::test {
2424

2525
TEST(ArrowUtilsTest, TestCreateProjection) {
26-
std::vector<DataField> read_fields = {DataField(1, arrow::field("k1", arrow::int32())),
27-
DataField(3, arrow::field("p1", arrow::int32())),
28-
DataField(5, arrow::field("s1", arrow::utf8())),
29-
DataField(6, arrow::field("v0", arrow::float64())),
30-
DataField(7, arrow::field("v1", arrow::boolean()))};
31-
auto read_schema = DataField::ConvertDataFieldsToArrowSchema(read_fields);
26+
arrow::FieldVector file_fields = {
27+
arrow::field("k0", arrow::int32()), arrow::field("k1", arrow::int32()),
28+
arrow::field("p1", arrow::int32()), arrow::field("s1", arrow::utf8()),
29+
arrow::field("v0", arrow::float64()), arrow::field("v1", arrow::boolean()),
30+
arrow::field("s0", arrow::utf8())};
31+
auto file_schema = arrow::schema(file_fields);
3232

33-
std::vector<DataField> file_fields = {DataField(0, arrow::field("k0", arrow::int32())),
34-
DataField(1, arrow::field("k1", arrow::int32())),
35-
DataField(3, arrow::field("p1", arrow::int32())),
36-
DataField(5, arrow::field("s1", arrow::utf8())),
37-
DataField(6, arrow::field("v0", arrow::float64())),
38-
DataField(7, arrow::field("v1", arrow::boolean())),
39-
DataField(4, arrow::field("s0", arrow::utf8()))};
40-
auto file_schema = DataField::ConvertDataFieldsToArrowSchema(file_fields);
41-
42-
auto projection = ArrowUtils::CreateProjection(file_schema, read_schema->fields());
43-
std::vector<int32_t> expected_projection = {1, 2, 3, 4, 5};
44-
ASSERT_EQ(projection, expected_projection);
33+
{
34+
// normal case
35+
arrow::FieldVector read_fields = {
36+
arrow::field("k1", arrow::int32()), arrow::field("p1", arrow::int32()),
37+
arrow::field("s1", arrow::utf8()), arrow::field("v0", arrow::float64()),
38+
arrow::field("v1", arrow::boolean())};
39+
auto read_schema = arrow::schema(read_fields);
40+
ASSERT_OK_AND_ASSIGN(std::vector<int32_t> projection,
41+
ArrowUtils::CreateProjection(file_schema, read_schema->fields()));
42+
std::vector<int32_t> expected_projection = {1, 2, 3, 4, 5};
43+
ASSERT_EQ(projection, expected_projection);
44+
}
45+
{
46+
// duplicate read field
47+
arrow::FieldVector read_fields = {
48+
arrow::field("k1", arrow::int32()), arrow::field("p1", arrow::int32()),
49+
arrow::field("s1", arrow::utf8()), arrow::field("v0", arrow::float64()),
50+
arrow::field("v0", arrow::float64()), arrow::field("v1", arrow::boolean())};
51+
auto read_schema = arrow::schema(read_fields);
52+
ASSERT_OK_AND_ASSIGN(std::vector<int32_t> projection,
53+
ArrowUtils::CreateProjection(file_schema, read_schema->fields()));
54+
std::vector<int32_t> expected_projection = {1, 2, 3, 4, 4, 5};
55+
ASSERT_EQ(projection, expected_projection);
56+
}
57+
{
58+
// duplicate read field, and sizeof(read_fields) > sizeof(file_fields)
59+
arrow::FieldVector read_fields = {
60+
arrow::field("k1", arrow::int32()), arrow::field("p1", arrow::int32()),
61+
arrow::field("s1", arrow::utf8()), arrow::field("v0", arrow::float64()),
62+
arrow::field("v0", arrow::float64()), arrow::field("v0", arrow::float64()),
63+
arrow::field("v0", arrow::float64()), arrow::field("v0", arrow::float64()),
64+
arrow::field("v1", arrow::boolean())};
65+
auto read_schema = arrow::schema(read_fields);
66+
ASSERT_OK_AND_ASSIGN(std::vector<int32_t> projection,
67+
ArrowUtils::CreateProjection(file_schema, read_schema->fields()));
68+
std::vector<int32_t> expected_projection = {1, 2, 3, 4, 4, 4, 4, 4, 5};
69+
ASSERT_EQ(projection, expected_projection);
70+
}
71+
{
72+
// read field not found in file schema
73+
arrow::FieldVector read_fields = {
74+
arrow::field("k1", arrow::int32()), arrow::field("p1", arrow::int32()),
75+
arrow::field("s1", arrow::utf8()), arrow::field("v2", arrow::float64()),
76+
arrow::field("v1", arrow::boolean())};
77+
auto read_schema = arrow::schema(read_fields);
78+
ASSERT_NOK_WITH_MSG(ArrowUtils::CreateProjection(file_schema, read_schema->fields()),
79+
"Field 'v2' not found or duplicate in file schema");
80+
}
81+
{
82+
// duplicate field in file schema
83+
arrow::FieldVector file_fields_dup = {
84+
arrow::field("k0", arrow::int32()), arrow::field("k1", arrow::int32()),
85+
arrow::field("p1", arrow::int32()), arrow::field("s1", arrow::utf8()),
86+
arrow::field("v0", arrow::float64()), arrow::field("v1", arrow::boolean()),
87+
arrow::field("v1", arrow::boolean()), arrow::field("s0", arrow::utf8())};
88+
auto file_schema_dup = arrow::schema(file_fields_dup);
89+
arrow::FieldVector read_fields = {
90+
arrow::field("k1", arrow::int32()), arrow::field("p1", arrow::int32()),
91+
arrow::field("s1", arrow::utf8()), arrow::field("v1", arrow::float64()),
92+
arrow::field("v1", arrow::boolean())};
93+
auto read_schema = arrow::schema(read_fields);
94+
ASSERT_NOK_WITH_MSG(ArrowUtils::CreateProjection(file_schema_dup, read_schema->fields()),
95+
"Field 'v1' not found or duplicate in file schema");
96+
}
97+
{
98+
arrow::FieldVector read_fields = {
99+
arrow::field("k1", arrow::int32()), arrow::field("p1", arrow::int32()),
100+
arrow::field("s1", arrow::utf8()), arrow::field("v0", arrow::float64()),
101+
arrow::field("v1", arrow::boolean())};
102+
auto read_schema = arrow::schema(read_fields);
103+
ASSERT_OK_AND_ASSIGN(std::vector<int32_t> projection,
104+
ArrowUtils::CreateProjection(file_schema, read_schema->fields()));
105+
std::vector<int32_t> expected_projection = {1, 2, 3, 4, 5};
106+
ASSERT_EQ(projection, expected_projection);
107+
}
45108
}
46109

47110
} // namespace paimon::test

src/paimon/core/operation/merge_file_split_read.cpp

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -106,8 +106,9 @@ Result<std::unique_ptr<MergeFileSplitRead>> MergeFileSplitRead::Create(
106106
int32_t key_arity = trimmed_primary_key.size();
107107

108108
// projection is the mapping from value_schema in KeyValue object to raw_read_schema
109-
std::vector<int32_t> projection =
110-
ArrowUtils::CreateProjection(value_schema, context->GetReadSchema()->fields());
109+
PAIMON_ASSIGN_OR_RAISE(
110+
std::vector<int32_t> projection,
111+
ArrowUtils::CreateProjection(value_schema, context->GetReadSchema()->fields()));
111112

112113
return std::unique_ptr<MergeFileSplitRead>(new MergeFileSplitRead(
113114
path_factory, context,

src/paimon/format/avro/avro_file_batch_reader.cpp

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -159,7 +159,8 @@ Status AvroFileBatchReader::SetReadSchema(::ArrowSchema* read_schema,
159159
Result<std::set<size_t>> AvroFileBatchReader::CalculateReadFieldsProjection(
160160
const std::shared_ptr<::arrow::Schema>& file_schema, const arrow::FieldVector& read_fields) {
161161
std::set<size_t> projection_set;
162-
auto projection = ArrowUtils::CreateProjection(file_schema, read_fields);
162+
PAIMON_ASSIGN_OR_RAISE(std::vector<int32_t> projection,
163+
ArrowUtils::CreateProjection(file_schema, read_fields));
163164
int32_t prev_index = -1;
164165
for (auto& index : projection) {
165166
if (index <= prev_index) {

src/paimon/format/avro/avro_format_writer_test.cpp

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -118,7 +118,6 @@ class AvroFormatWriterTest : public ::testing::Test {
118118
ASSERT_OK_AND_ASSIGN(std::shared_ptr<InputStream> input_stream, fs_->Open(file_path));
119119
ASSERT_OK_AND_ASSIGN(auto file_reader,
120120
AvroFileBatchReader::Create(input_stream, 1024, pool_));
121-
auto& reader = file_reader->reader_;
122121
ASSERT_OK_AND_ASSIGN(uint64_t num_rows, file_reader->GetNumberOfRows());
123122
ASSERT_EQ(num_rows, row_count);
124123

src/paimon/format/avro/avro_stats_extractor_test.cpp

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -120,12 +120,16 @@ TEST_F(AvroStatsExtractorTest, TestPrimitiveStatsExtractor) {
120120
ASSERT_TRUE(arrow::ExportSchema(*schema, &arrow_schema).ok());
121121
ASSERT_OK_AND_ASSIGN(auto extractor, format.CreateStatsExtractor(&arrow_schema));
122122
auto fs = std::make_shared<LocalFileSystem>();
123-
ASSERT_OK_AND_ASSIGN(auto results, extractor->Extract(fs, file_path, GetDefaultPool()));
123+
ASSERT_OK_AND_ASSIGN(auto stats_with_info,
124+
extractor->ExtractWithFileInfo(fs, file_path, GetDefaultPool()));
125+
const auto& column_stats = stats_with_info.first;
126+
const auto& file_stats = stats_with_info.second;
124127

125-
ASSERT_EQ(results.size(), 20);
126-
for (const auto& stats : results) {
128+
ASSERT_EQ(column_stats.size(), 20);
129+
for (const auto& stats : column_stats) {
127130
ASSERT_EQ(stats->ToString(), "min null, max null, null count null");
128131
}
132+
ASSERT_EQ(3, stats_with_info.second.GetRowCount());
129133
}
130134

131135
TEST_F(AvroStatsExtractorTest, TestNestedType) {

0 commit comments

Comments
 (0)