Skip to content

Commit c68be86

Browse files
authored
feat(blob): Support blob-descriptor-field for inline blob descriptor storage (alibaba#300)
1 parent 96d21b3 commit c68be86

47 files changed

Lines changed: 2571 additions & 419 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

include/paimon/data/blob.h

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -97,7 +97,8 @@ class PAIMON_EXPORT Blob {
9797
/// @param metadata A map of key-value metadata to be attached to the field.
9898
/// @return A result containing a unique pointer to the generated `ArrowSchema` or an error.
9999
static Result<std::unique_ptr<::ArrowSchema>> ArrowField(
100-
const std::string& field_name, std::unordered_map<std::string, std::string> metadata = {});
100+
const std::string& field_name, bool nullable = false,
101+
std::unordered_map<std::string, std::string> metadata = {});
101102

102103
private:
103104
class Impl;

include/paimon/defs.h

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -365,7 +365,7 @@ struct PAIMON_EXPORT Options {
365365
/// "partition.legacy-name" - The legacy partition name is using `ToString` for all types. If
366366
/// false, using casting to string for all types. Default value is "true".
367367
static const char PARTITION_GENERATE_LEGACY_NAME[];
368-
/// "blob-as-descriptor" - Read and write blob field using blob descriptor rather than blob
368+
/// "blob-as-descriptor" - Read blob field using blob descriptor rather than blob
369369
/// bytes. Default value is "false".
370370
static const char BLOB_AS_DESCRIPTOR[];
371371
/// "blob-field" - Specifies column names that should be stored as blob type. This is used

src/paimon/CMakeLists.txt

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -159,6 +159,7 @@ set(PAIMON_CORE_SRCS
159159
core/bucket/hive_bucket_function.cpp
160160
core/bucket/mod_bucket_function.cpp
161161
core/bucket/bucket_id_calculator.cpp
162+
core/casting/binary_to_blob_cast_executor.cpp
162163
core/casting/binary_to_string_cast_executor.cpp
163164
core/casting/boolean_to_decimal_cast_executor.cpp
164165
core/casting/boolean_to_numeric_cast_executor.cpp
@@ -222,6 +223,7 @@ set(PAIMON_CORE_SRCS
222223
core/io/key_value_meta_projection_consumer.cpp
223224
core/io/key_value_projection_consumer.cpp
224225
core/io/key_value_projection_reader.cpp
226+
core/io/external_storage_blob_writer.cpp
225227
core/io/multiple_blob_file_writer.cpp
226228
core/io/rolling_blob_file_writer.cpp
227229
core/manifest/file_kind.cpp
@@ -601,6 +603,7 @@ if(PAIMON_BUILD_TESTS)
601603
core/io/file_index_evaluator_test.cpp
602604
core/io/single_file_writer_test.cpp
603605
core/io/rolling_blob_file_writer_test.cpp
606+
core/io/external_storage_blob_writer_test.cpp
604607
core/global_index/indexed_split_test.cpp
605608
core/manifest/file_source_test.cpp
606609
core/manifest/file_kind_test.cpp

src/paimon/common/data/blob.cpp

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -105,8 +105,9 @@ Result<PAIMON_UNIQUE_PTR<Bytes>> Blob::ToData(const std::shared_ptr<FileSystem>&
105105
}
106106

107107
Result<std::unique_ptr<ArrowSchema>> Blob::ArrowField(
108-
const std::string& field_name, std::unordered_map<std::string, std::string> metadata) {
109-
auto blob_field = BlobUtils::ToArrowField(field_name, /*nullable=*/false, metadata);
108+
const std::string& field_name, bool nullable,
109+
std::unordered_map<std::string, std::string> metadata) {
110+
auto blob_field = BlobUtils::ToArrowField(field_name, nullable, metadata);
110111
auto field = std::make_unique<::ArrowSchema>();
111112
PAIMON_RETURN_NOT_OK_FROM_ARROW(arrow::ExportField(*blob_field, field.get()));
112113
return field;

src/paimon/common/data/blob_descriptor.h

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -38,7 +38,7 @@ namespace paimon {
3838
/// | 13 + N | offset | long | 8 |
3939
/// | 21 + N | length | long | 8 |
4040

41-
class BlobDescriptor {
41+
class PAIMON_EXPORT BlobDescriptor {
4242
public:
4343
static Result<std::unique_ptr<BlobDescriptor>> Create(const std::string& uri, int64_t offset,
4444
int64_t length);

src/paimon/common/data/blob_test.cpp

Lines changed: 5 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -144,38 +144,34 @@ TEST_F(BlobTest, TestNewInputStreamWithDynamicLength) {
144144
}
145145

146146
TEST_F(BlobTest, TestArrowField) {
147-
{
148-
// basic: field name, non-nullable by default
149-
ASSERT_OK_AND_ASSIGN(auto schema, Blob::ArrowField("my_blob"));
147+
for (bool nullable : {false, true}) {
148+
ASSERT_OK_AND_ASSIGN(auto schema, Blob::ArrowField("my_blob", nullable));
150149
ASSERT_NE(schema, nullptr);
151150

152-
// import back to arrow::Field to verify
153151
auto field_result = arrow::ImportField(schema.get());
154152
ASSERT_TRUE(field_result.ok());
155153
auto field = field_result.ValueUnsafe();
156154

157155
ASSERT_EQ(field->name(), "my_blob");
158156
ASSERT_EQ(field->type()->id(), arrow::Type::LARGE_BINARY);
159-
ASSERT_FALSE(field->nullable());
157+
ASSERT_EQ(field->nullable(), nullable);
160158
ASSERT_TRUE(field->HasMetadata());
161159
auto extension_type = field->metadata()->Get("paimon.extension.type");
162160
ASSERT_TRUE(extension_type.ok());
163161
ASSERT_EQ(extension_type.ValueUnsafe(), "paimon.type.blob");
164162
}
165163
{
166-
// with custom metadata
167164
std::unordered_map<std::string, std::string> custom_metadata = {
168165
{"custom_key", "custom_value"}};
169-
ASSERT_OK_AND_ASSIGN(auto schema, Blob::ArrowField("meta_blob", custom_metadata));
166+
ASSERT_OK_AND_ASSIGN(auto schema,
167+
Blob::ArrowField("meta_blob", /*nullable=*/false, custom_metadata));
170168
auto field = arrow::ImportField(schema.get()).ValueUnsafe();
171169
ASSERT_EQ(field->name(), "meta_blob");
172170
ASSERT_FALSE(field->nullable());
173171
ASSERT_TRUE(field->HasMetadata());
174-
// blob extension metadata should be present
175172
auto extension_type = field->metadata()->Get("paimon.extension.type");
176173
ASSERT_TRUE(extension_type.ok());
177174
ASSERT_EQ(extension_type.ValueUnsafe(), "paimon.type.blob");
178-
// custom metadata should also be present
179175
auto custom_val = field->metadata()->Get("custom_key");
180176
ASSERT_TRUE(custom_val.ok());
181177
ASSERT_EQ(custom_val.ValueUnsafe(), "custom_value");

src/paimon/common/data/blob_utils.cpp

Lines changed: 92 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -17,65 +17,77 @@
1717
#include "paimon/common/data/blob_utils.h"
1818

1919
#include <cstddef>
20-
#include <utility>
20+
#include <set>
2121
#include <vector>
2222

2323
#include "arrow/api.h"
2424
#include "arrow/array/array_nested.h"
2525
#include "arrow/type.h"
26+
#include "fmt/format.h"
2627
#include "paimon/common/data/blob_defs.h"
28+
#include "paimon/common/data/blob_descriptor.h"
29+
#include "paimon/common/types/data_field.h"
2730
#include "paimon/common/utils/arrow/status_utils.h"
2831
#include "paimon/common/utils/string_utils.h"
29-
3032
namespace arrow {
3133
class Array;
3234
}
3335

3436
namespace paimon {
35-
3637
BlobUtils::SeparatedSchemas BlobUtils::SeparateBlobSchema(
37-
const std::shared_ptr<arrow::Schema>& schema) {
38-
std::vector<std::shared_ptr<arrow::Field>> remaining_fields;
38+
const std::shared_ptr<arrow::Schema>& schema, const std::set<std::string>& inline_fields) {
39+
std::vector<std::shared_ptr<arrow::Field>> main_fields;
3940
std::vector<std::shared_ptr<arrow::Field>> blob_fields;
40-
for (auto i = 0; i < schema->num_fields(); i++) {
41+
for (int32_t i = 0; i < schema->num_fields(); i++) {
4142
auto field = schema->field(i);
42-
if (IsBlobField(field)) {
43+
if (IsBlobField(field) && inline_fields.count(field->name()) == 0) {
44+
// Non-inline BLOB -> goes to blob file
4345
blob_fields.emplace_back(field);
4446
} else {
45-
remaining_fields.emplace_back(field);
47+
// Non-blob fields OR inline BLOB fields -> stay in main
48+
main_fields.emplace_back(field);
4649
}
4750
}
4851
SeparatedSchemas result;
49-
result.main_schema = arrow::schema(remaining_fields);
52+
result.main_schema = arrow::schema(main_fields);
5053
result.blob_schema = arrow::schema(blob_fields);
5154
return result;
5255
}
5356

5457
Result<BlobUtils::SeparatedStructArrays> BlobUtils::SeparateBlobArray(
55-
const std::shared_ptr<arrow::StructArray>& struct_array) {
58+
const std::shared_ptr<arrow::StructArray>& struct_array,
59+
const std::set<std::string>& inline_fields) {
5660
std::shared_ptr<arrow::StructType> old_type =
5761
std::static_pointer_cast<arrow::StructType>(struct_array->type());
5862
const auto& old_fields = old_type->fields();
5963
const auto& old_arrays = struct_array->fields();
6064

61-
std::vector<std::shared_ptr<arrow::Field>> remaining_fields;
62-
std::vector<std::shared_ptr<arrow::Array>> remaining_arrays;
63-
std::vector<std::shared_ptr<arrow::Field>> blob_fields;
64-
std::vector<std::shared_ptr<arrow::Array>> blob_arrays;
65+
arrow::ArrayVector main_arrays;
66+
arrow::ArrayVector blob_arrays;
67+
arrow::FieldVector main_fields;
68+
arrow::FieldVector blob_fields;
6569

6670
for (size_t i = 0; i < old_fields.size(); i++) {
67-
if (IsBlobField(old_fields[i])) {
71+
if (IsBlobField(old_fields[i]) && inline_fields.count(old_fields[i]->name()) == 0) {
6872
blob_fields.push_back(old_fields[i]);
6973
blob_arrays.push_back(old_arrays[i]);
7074
} else {
71-
remaining_fields.push_back(old_fields[i]);
72-
remaining_arrays.push_back(old_arrays[i]);
75+
main_fields.push_back(old_fields[i]);
76+
main_arrays.push_back(old_arrays[i]);
7377
}
7478
}
7579

80+
if (blob_fields.empty()) {
81+
return Status::Invalid(
82+
"SeparateBlobArray expects at least one non-inline blob field, but got none.");
83+
}
84+
if (main_fields.empty()) {
85+
return Status::Invalid("SeparateBlobArray expects at least one main field, but got none.");
86+
}
87+
7688
SeparatedStructArrays result;
7789
PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(result.main_array,
78-
arrow::StructArray::Make(remaining_arrays, remaining_fields));
90+
arrow::StructArray::Make(main_arrays, main_fields));
7991
PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(result.blob_array,
8092
arrow::StructArray::Make(blob_arrays, blob_fields));
8193
return result;
@@ -114,4 +126,66 @@ std::shared_ptr<arrow::Field> BlobUtils::ToArrowField(
114126
return arrow::field(field_name, arrow::large_binary(), nullable,
115127
std::make_shared<arrow::KeyValueMetadata>(metadata));
116128
}
129+
130+
Status BlobUtils::ValidateInlineBlobDescriptors(
131+
const std::shared_ptr<arrow::StructArray>& struct_array,
132+
const std::set<std::string>& inline_descriptor_fields) {
133+
if (inline_descriptor_fields.empty()) {
134+
return Status::OK();
135+
}
136+
if (!struct_array) {
137+
return Status::Invalid("array in ValidateInlineBlobDescriptors must be a struct_array");
138+
}
139+
for (const auto& field_name : inline_descriptor_fields) {
140+
auto field_array = struct_array->GetFieldByName(field_name);
141+
if (!field_array) {
142+
continue;
143+
}
144+
const auto* binary_array =
145+
arrow::internal::checked_cast<const arrow::LargeBinaryArray*>(field_array.get());
146+
if (!binary_array) {
147+
return Status::Invalid(
148+
fmt::format("cannot cast array for field {} to LargeBinaryArray", field_name));
149+
}
150+
for (int64_t row = 0; row < binary_array->length(); ++row) {
151+
if (binary_array->IsNull(row)) {
152+
continue;
153+
}
154+
auto value = binary_array->GetView(row);
155+
PAIMON_ASSIGN_OR_RAISE(bool is_descriptor,
156+
BlobDescriptor::IsBlobDescriptor(value.data(), value.size()));
157+
if (!is_descriptor) {
158+
return Status::Invalid(fmt::format(
159+
"BLOB inline field {} configured by blob-descriptor-field or blob-view-field "
160+
"require values to be a BlobDescriptor or BlobViewStruct.",
161+
field_name));
162+
}
163+
}
164+
}
165+
return Status::OK();
166+
}
167+
168+
std::vector<DataField> BlobUtils::ConvertBlobInlineDataFields(
169+
const std::vector<DataField>& data_fields, const std::vector<std::string>& blob_inline_fields) {
170+
if (blob_inline_fields.empty()) {
171+
return data_fields;
172+
}
173+
174+
std::set<std::string> blob_inline_field_set(blob_inline_fields.begin(),
175+
blob_inline_fields.end());
176+
std::vector<DataField> converted_fields;
177+
converted_fields.reserve(data_fields.size());
178+
for (const auto& data_field : data_fields) {
179+
if (blob_inline_field_set.find(data_field.Name()) == blob_inline_field_set.end()) {
180+
converted_fields.push_back(data_field);
181+
continue;
182+
}
183+
184+
auto binary_field = arrow::field(data_field.Name(), arrow::binary(), data_field.Nullable(),
185+
data_field.ArrowField()->metadata());
186+
converted_fields.emplace_back(data_field.Id(), binary_field, data_field.Description());
187+
}
188+
return converted_fields;
189+
}
190+
117191
} // namespace paimon

src/paimon/common/data/blob_utils.h

Lines changed: 30 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -17,8 +17,10 @@
1717
#pragma once
1818

1919
#include <memory>
20+
#include <set>
2021
#include <string>
2122
#include <unordered_map>
23+
#include <vector>
2224

2325
#include "paimon/result.h"
2426
#include "paimon/visibility.h"
@@ -30,6 +32,10 @@ class Schema;
3032
class StructArray;
3133
} // namespace arrow
3234

35+
namespace paimon {
36+
class DataField;
37+
} // namespace paimon
38+
3339
namespace paimon {
3440
/// Utils for blob type.
3541
class PAIMON_EXPORT BlobUtils {
@@ -38,23 +44,29 @@ class PAIMON_EXPORT BlobUtils {
3844
~BlobUtils() = delete;
3945

4046
struct SeparatedSchemas {
41-
/// Non-blob fields
47+
/// Non-blob fields (includes inline blob fields when inline_fields is provided)
4248
std::shared_ptr<arrow::Schema> main_schema;
43-
/// Blob fields only
49+
/// Blob fields that go to separate .blob files
4450
std::shared_ptr<arrow::Schema> blob_schema;
4551
};
4652

4753
struct SeparatedStructArrays {
48-
/// Non-blob fields
54+
/// Non-blob fields (includes inline blob fields when inline_fields is provided)
4955
std::shared_ptr<arrow::StructArray> main_array;
50-
/// Blob fields only
56+
/// Blob fields that go to separate .blob files
5157
std::shared_ptr<arrow::StructArray> blob_array;
5258
};
5359

54-
static SeparatedSchemas SeparateBlobSchema(const std::shared_ptr<arrow::Schema>& schema);
60+
/// Separates schema with inline field awareness.
61+
/// BLOB fields in inline_fields stay in main_schema; others go to blob_schema.
62+
static SeparatedSchemas SeparateBlobSchema(const std::shared_ptr<arrow::Schema>& schema,
63+
const std::set<std::string>& inline_fields);
5564

65+
/// Separates array with inline field awareness.
66+
/// BLOB fields in inline_fields stay in main_array; others go to blob_array.
5667
static Result<SeparatedStructArrays> SeparateBlobArray(
57-
const std::shared_ptr<arrow::StructArray>& struct_array);
68+
const std::shared_ptr<arrow::StructArray>& struct_array,
69+
const std::set<std::string>& inline_fields);
5870

5971
static bool IsBlobField(const std::shared_ptr<arrow::Field>& field);
6072
static bool IsBlobMetadata(const std::shared_ptr<const arrow::KeyValueMetadata>& metadata);
@@ -63,6 +75,18 @@ class PAIMON_EXPORT BlobUtils {
6375
static std::shared_ptr<arrow::Field> ToArrowField(
6476
const std::string& field_name, bool nullable = false,
6577
std::unordered_map<std::string, std::string> metadata = {});
78+
79+
static Status ValidateInlineBlobDescriptors(
80+
const std::shared_ptr<arrow::StructArray>& struct_array,
81+
const std::set<std::string>& inline_descriptor_fields);
82+
83+
/// Converts inline blob DataFields from large_binary to binary type.
84+
/// Inline blob fields use large_binary in the table schema (because they are BLOB type),
85+
/// but are stored as binary in data files. This conversion aligns the field type with
86+
/// the actual on-disk storage format for correct reading.
87+
static std::vector<DataField> ConvertBlobInlineDataFields(
88+
const std::vector<DataField>& data_fields,
89+
const std::vector<std::string>& blob_inline_fields);
6690
};
6791

6892
} // namespace paimon

0 commit comments

Comments
 (0)