github-actions[bot] commented on code in PR #67325:
URL: https://github.com/apache/doris/pull/67325#discussion_r3976194119
##########
be/src/service/internal_service.cpp:
##########
@@ -923,7 +923,11 @@ void
PInternalService::fetch_table_schema(google::protobuf::RpcController* contr
for (const auto& col_type : col_types) {
DORIS_CHECK(col_type != nullptr);
PTypeDesc* type_desc = result->add_column_types();
- if (col_type->get_primitive_type() == INVALID_TYPE) {
+ if (col_type->is_null_literal()) {
Review Comment:
[P2] Keep schema discovery and execution on the same BE capability set.
During Doris's [documented BE-first rolling
upgrade](https://doris.apache.org/docs/dev/admin-manual/cluster-management/upgrade/),
`local(..., "shared_storage"="true", "format"="lance")` picks a random
`backendIdForRequest` for this RPC, but `TVFScanNode` leaves execution on every
BE when the public `backendId` is -1. A new BE can therefore expose these
ordinary `NULL_TYPE`/`FLOAT`/`JSONB`/`BIGINT` nodes to an old FE, which
generically accepts them, and the one Lance split can then run on an old BE.
BFloat16 is a concrete failure: the old reader sends FixedSizeBinary(2)
directly to the FLOAT SerDe, which rejects the two-byte buffer under default
Arrow validation. Please pin execution to the schema BE or carry/filter an
explicit BE capability until all eligible BEs support the new materialization,
and add a mixed-BE route test.
##########
be/src/format_v2/lance/lance_reader_helper.cpp:
##########
@@ -178,8 +272,333 @@ Status arrow_field_to_doris_type(const
std::shared_ptr<arrow::Field>& field,
}
}
+// Determine whether a field subtree contains values that require Lance
normalization.
+Status field_requires_lance_normalization(const std::shared_ptr<arrow::Field>&
field,
+ bool* requires_normalization) {
+ DORIS_CHECK(field != nullptr);
+ DORIS_CHECK(requires_normalization != nullptr);
+
+ LanceExtensionKind extension_kind;
+ std::shared_ptr<arrow::DataType> storage_type;
+ RETURN_IF_ERROR(get_lance_extension(field, &extension_kind,
&storage_type));
+ bool required = extension_kind == LanceExtensionKind::BFLOAT16 ||
+ field->type()->id() == arrow::Type::EXTENSION;
+ for (const auto& child : storage_type->fields()) {
+ bool child_required = false;
+ RETURN_IF_ERROR(field_requires_lance_normalization(child,
&child_required));
+ required |= child_required;
+ }
+ *requires_normalization = required;
+ return Status::OK();
+}
+
+// Widen little-endian Lance BFloat16 values to Arrow Float32 without
precision loss.
+Status convert_bfloat16_array(const std::shared_ptr<arrow::Array>& array,
+ std::shared_ptr<arrow::Array>* normalized) {
+ DORIS_CHECK(array != nullptr);
+ DORIS_CHECK(normalized != nullptr);
+ const auto fixed_binary =
std::dynamic_pointer_cast<arrow::FixedSizeBinaryArray>(array);
+ if (fixed_binary == nullptr || fixed_binary->byte_width() != 2) {
+ return Status::InvalidArgument("invalid Lance BFloat16 array storage:
{}",
+ array->type()->ToString());
+ }
+ if (config::enable_arrow_input_validation) {
+ check_arrow_fixed_width_buffer(*fixed_binary, sizeof(uint16_t));
+ }
+
+ arrow::FloatBuilder builder;
+ auto arrow_status = builder.Reserve(fixed_binary->length());
+ if (!arrow_status.ok()) {
+ return Status::InternalError("reserve Lance BFloat16 output failed:
{}",
+ arrow_status.message());
+ }
+ for (int64_t row = 0; row < fixed_binary->length(); ++row) {
+ if (fixed_binary->IsNull(row)) {
+ arrow_status = builder.AppendNull();
+ } else {
+ const auto bits =
LittleEndian::Load16(fixed_binary->GetValue(row));
+ arrow_status =
builder.Append(std::bit_cast<float>(static_cast<uint32_t>(bits) << 16));
+ }
+ if (!arrow_status.ok()) {
+ return Status::InternalError("append Lance BFloat16 value failed:
{}",
+ arrow_status.message());
+ }
+ }
+ std::shared_ptr<arrow::FloatArray> result;
+ arrow_status = builder.Finish(&result);
+ if (!arrow_status.ok()) {
+ return Status::InternalError("finish Lance BFloat16 conversion failed:
{}",
+ arrow_status.message());
+ }
+ *normalized = std::move(result);
+ return Status::OK();
+}
+
+// Copy parent metadata into an offset-zero view with a matching validity
bitmap.
+Status rebase_lance_parent_data(const std::shared_ptr<arrow::Array>& array,
+ std::shared_ptr<arrow::ArrayData>*
rebased_data) {
+ auto data = array->data()->Copy();
+ data->offset = 0;
+ data->SetNullCount(array->null_count());
+ const auto& null_bitmap = array->data()->buffers[0];
+ if (null_bitmap != nullptr) {
+ auto bitmap_result =
+ arrow::internal::CopyBitmap(arrow::default_memory_pool(),
null_bitmap->data(),
+ array->offset(), array->length());
+ if (!bitmap_result.ok()) {
+ return Status::InternalError("copy sliced Lance validity bitmap
failed: {}",
+ bitmap_result.status().message());
+ }
+ data->buffers[0] = std::move(bitmap_result).ValueUnsafe();
+ }
+ *rebased_data = std::move(data);
+ return Status::OK();
+}
+
+// Compact a sliced variable-offset parent and its child to the visible value
interval.
+template <typename ArrayType, typename OffsetType>
+Status compact_lance_offset_array(const std::shared_ptr<arrow::Array>& array,
+ std::shared_ptr<arrow::Array>* compacted) {
+ const auto offset_array = std::dynamic_pointer_cast<ArrayType>(array);
+ if (offset_array == nullptr) {
+ return Status::InvalidArgument("invalid sliced Lance offset array: {}",
+ array->type()->ToString());
+ }
+ const auto child_begin =
static_cast<int64_t>(offset_array->value_offset(0));
+ const auto child_end =
static_cast<int64_t>(offset_array->value_offset(offset_array->length()));
+ const auto& values = offset_array->values();
+ if (child_begin < 0 || child_end < child_begin || child_end >
values->length()) {
+ return Status::InvalidArgument("invalid sliced Lance offsets [{}, {})
for child length {}",
+ child_begin, child_end,
values->length());
+ }
+ if (array->offset() == 0 && child_begin == 0 && child_end ==
values->length()) {
+ *compacted = array;
+ return Status::OK();
+ }
+
+ arrow::TypedBufferBuilder<OffsetType> offsets_builder;
+ auto arrow_status = offsets_builder.Reserve(array->length() + 1);
+ if (!arrow_status.ok()) {
+ return Status::InternalError("reserve sliced Lance offsets failed: {}",
+ arrow_status.message());
+ }
+ for (int64_t index = 0; index <= array->length(); ++index) {
+ arrow_status = offsets_builder.Append(
+ static_cast<OffsetType>(offset_array->value_offset(index) -
child_begin));
+ if (!arrow_status.ok()) {
+ return Status::InternalError("append sliced Lance offset failed:
{}",
+ arrow_status.message());
+ }
+ }
+ std::shared_ptr<arrow::Buffer> offsets;
+ arrow_status = offsets_builder.Finish(&offsets);
+ if (!arrow_status.ok()) {
+ return Status::InternalError("finish sliced Lance offsets failed: {}",
+ arrow_status.message());
+ }
+
+ std::shared_ptr<arrow::ArrayData> rebased_data;
+ RETURN_IF_ERROR(rebase_lance_parent_data(array, &rebased_data));
+ rebased_data->buffers[1] = std::move(offsets);
+ rebased_data->child_data[0] = values->Slice(child_begin, child_end -
child_begin)->data();
+ *compacted = arrow::MakeArray(std::move(rebased_data));
+ return Status::OK();
+}
+
+// Compact sliced nested parents so recursive normalization sees only visible
child values.
+Status compact_lance_nested_array(const std::shared_ptr<arrow::Array>& array,
+ std::shared_ptr<arrow::Array>* compacted) {
+ switch (array->type_id()) {
+ case arrow::Type::LIST:
+ return compact_lance_offset_array<arrow::ListArray, int32_t>(array,
compacted);
+ case arrow::Type::LARGE_LIST:
+ return compact_lance_offset_array<arrow::LargeListArray,
int64_t>(array, compacted);
+ case arrow::Type::MAP:
+ return compact_lance_offset_array<arrow::MapArray, int32_t>(array,
compacted);
+ case arrow::Type::FIXED_SIZE_LIST: {
+ const auto list =
std::dynamic_pointer_cast<arrow::FixedSizeListArray>(array);
+ if (list == nullptr) {
+ return Status::InvalidArgument("invalid sliced Lance fixed-size
list array: {}",
+ array->type()->ToString());
+ }
+ const auto child_begin = list->value_offset(0);
+ const auto child_length = list->length() * list->value_length();
+ const auto& values = list->values();
+ if (child_begin < 0 || child_length < 0 || child_begin >
values->length() - child_length) {
+ return Status::InvalidArgument(
+ "invalid sliced Lance fixed-size list range [{}, {}) for
child length {}",
+ child_begin, child_begin + child_length, values->length());
+ }
+ if (array->offset() == 0 && child_begin == 0 && child_length ==
values->length()) {
+ *compacted = array;
+ return Status::OK();
+ }
+ std::shared_ptr<arrow::ArrayData> rebased_data;
+ RETURN_IF_ERROR(rebase_lance_parent_data(array, &rebased_data));
+ rebased_data->child_data[0] = values->Slice(child_begin,
child_length)->data();
+ *compacted = arrow::MakeArray(std::move(rebased_data));
+ return Status::OK();
+ }
+ case arrow::Type::STRUCT: {
+ const auto struct_array =
std::dynamic_pointer_cast<arrow::StructArray>(array);
+ if (struct_array == nullptr) {
+ return Status::InvalidArgument("invalid sliced Lance struct array:
{}",
+ array->type()->ToString());
+ }
+ bool requires_compaction = array->offset() != 0;
+ for (const auto& child : array->data()->child_data) {
+ requires_compaction |= child->length != array->length();
+ }
+ if (!requires_compaction) {
+ *compacted = array;
+ return Status::OK();
+ }
+ std::shared_ptr<arrow::ArrayData> rebased_data;
+ RETURN_IF_ERROR(rebase_lance_parent_data(array, &rebased_data));
+ for (int child_idx = 0; child_idx <
static_cast<int>(struct_array->fields().size());
+ ++child_idx) {
+ rebased_data->child_data[child_idx] =
struct_array->field(child_idx)->data();
Review Comment:
[P2] Rebase every child when resetting a sliced Struct parent to offset
zero. [Arrow 24
`StructArray::field()`](https://github.com/apache/arrow/blob/apache-arrow-24.0.0/cpp/src/arrow/array/array_nested.cc#L992-L1008)
applies the parent's slice, so in the changed sliced-Map case the ordinary
String key installed here still has offset 2; only the BFloat16 item is rebuilt
at offset zero. `DataTypeMapSerDe` then asks the key SerDe for `[0,3)`, which
fails the default zero-offset validation (or reads k0/k1/k2 rather than
k2/k3/k4 with validation disabled). The test's direct `GetString()` assertions
are offset-aware and therefore miss the Doris failure. Please compact/rebase
unchanged siblings too, and materialize the mixed Map/Struct through
`_fill_block_from_record_batch`.
##########
be/src/format_v2/lance/lance_reader_helper.cpp:
##########
@@ -178,8 +272,333 @@ Status arrow_field_to_doris_type(const
std::shared_ptr<arrow::Field>& field,
}
}
+// Determine whether a field subtree contains values that require Lance
normalization.
+Status field_requires_lance_normalization(const std::shared_ptr<arrow::Field>&
field,
+ bool* requires_normalization) {
+ DORIS_CHECK(field != nullptr);
+ DORIS_CHECK(requires_normalization != nullptr);
+
+ LanceExtensionKind extension_kind;
+ std::shared_ptr<arrow::DataType> storage_type;
+ RETURN_IF_ERROR(get_lance_extension(field, &extension_kind,
&storage_type));
+ bool required = extension_kind == LanceExtensionKind::BFLOAT16 ||
+ field->type()->id() == arrow::Type::EXTENSION;
+ for (const auto& child : storage_type->fields()) {
+ bool child_required = false;
+ RETURN_IF_ERROR(field_requires_lance_normalization(child,
&child_required));
+ required |= child_required;
+ }
+ *requires_normalization = required;
+ return Status::OK();
+}
+
+// Widen little-endian Lance BFloat16 values to Arrow Float32 without
precision loss.
+Status convert_bfloat16_array(const std::shared_ptr<arrow::Array>& array,
+ std::shared_ptr<arrow::Array>* normalized) {
+ DORIS_CHECK(array != nullptr);
+ DORIS_CHECK(normalized != nullptr);
+ const auto fixed_binary =
std::dynamic_pointer_cast<arrow::FixedSizeBinaryArray>(array);
+ if (fixed_binary == nullptr || fixed_binary->byte_width() != 2) {
+ return Status::InvalidArgument("invalid Lance BFloat16 array storage:
{}",
+ array->type()->ToString());
+ }
+ if (config::enable_arrow_input_validation) {
+ check_arrow_fixed_width_buffer(*fixed_binary, sizeof(uint16_t));
+ }
+
+ arrow::FloatBuilder builder;
+ auto arrow_status = builder.Reserve(fixed_binary->length());
+ if (!arrow_status.ok()) {
+ return Status::InternalError("reserve Lance BFloat16 output failed:
{}",
+ arrow_status.message());
+ }
+ for (int64_t row = 0; row < fixed_binary->length(); ++row) {
+ if (fixed_binary->IsNull(row)) {
+ arrow_status = builder.AppendNull();
+ } else {
+ const auto bits =
LittleEndian::Load16(fixed_binary->GetValue(row));
+ arrow_status =
builder.Append(std::bit_cast<float>(static_cast<uint32_t>(bits) << 16));
+ }
+ if (!arrow_status.ok()) {
+ return Status::InternalError("append Lance BFloat16 value failed:
{}",
+ arrow_status.message());
+ }
+ }
+ std::shared_ptr<arrow::FloatArray> result;
+ arrow_status = builder.Finish(&result);
+ if (!arrow_status.ok()) {
+ return Status::InternalError("finish Lance BFloat16 conversion failed:
{}",
+ arrow_status.message());
+ }
+ *normalized = std::move(result);
+ return Status::OK();
+}
+
+// Copy parent metadata into an offset-zero view with a matching validity
bitmap.
+Status rebase_lance_parent_data(const std::shared_ptr<arrow::Array>& array,
+ std::shared_ptr<arrow::ArrayData>*
rebased_data) {
+ auto data = array->data()->Copy();
+ data->offset = 0;
+ data->SetNullCount(array->null_count());
+ const auto& null_bitmap = array->data()->buffers[0];
+ if (null_bitmap != nullptr) {
+ auto bitmap_result =
+ arrow::internal::CopyBitmap(arrow::default_memory_pool(),
null_bitmap->data(),
+ array->offset(), array->length());
+ if (!bitmap_result.ok()) {
+ return Status::InternalError("copy sliced Lance validity bitmap
failed: {}",
+ bitmap_result.status().message());
+ }
+ data->buffers[0] = std::move(bitmap_result).ValueUnsafe();
+ }
+ *rebased_data = std::move(data);
+ return Status::OK();
+}
+
+// Compact a sliced variable-offset parent and its child to the visible value
interval.
+template <typename ArrayType, typename OffsetType>
+Status compact_lance_offset_array(const std::shared_ptr<arrow::Array>& array,
+ std::shared_ptr<arrow::Array>* compacted) {
+ const auto offset_array = std::dynamic_pointer_cast<ArrayType>(array);
+ if (offset_array == nullptr) {
+ return Status::InvalidArgument("invalid sliced Lance offset array: {}",
+ array->type()->ToString());
+ }
+ const auto child_begin =
static_cast<int64_t>(offset_array->value_offset(0));
+ const auto child_end =
static_cast<int64_t>(offset_array->value_offset(offset_array->length()));
+ const auto& values = offset_array->values();
+ if (child_begin < 0 || child_end < child_begin || child_end >
values->length()) {
+ return Status::InvalidArgument("invalid sliced Lance offsets [{}, {})
for child length {}",
+ child_begin, child_end,
values->length());
+ }
+ if (array->offset() == 0 && child_begin == 0 && child_end ==
values->length()) {
+ *compacted = array;
+ return Status::OK();
+ }
+
+ arrow::TypedBufferBuilder<OffsetType> offsets_builder;
+ auto arrow_status = offsets_builder.Reserve(array->length() + 1);
+ if (!arrow_status.ok()) {
+ return Status::InternalError("reserve sliced Lance offsets failed: {}",
+ arrow_status.message());
+ }
+ for (int64_t index = 0; index <= array->length(); ++index) {
+ arrow_status = offsets_builder.Append(
+ static_cast<OffsetType>(offset_array->value_offset(index) -
child_begin));
+ if (!arrow_status.ok()) {
+ return Status::InternalError("append sliced Lance offset failed:
{}",
+ arrow_status.message());
+ }
+ }
+ std::shared_ptr<arrow::Buffer> offsets;
+ arrow_status = offsets_builder.Finish(&offsets);
+ if (!arrow_status.ok()) {
+ return Status::InternalError("finish sliced Lance offsets failed: {}",
+ arrow_status.message());
+ }
+
+ std::shared_ptr<arrow::ArrayData> rebased_data;
+ RETURN_IF_ERROR(rebase_lance_parent_data(array, &rebased_data));
+ rebased_data->buffers[1] = std::move(offsets);
+ rebased_data->child_data[0] = values->Slice(child_begin, child_end -
child_begin)->data();
+ *compacted = arrow::MakeArray(std::move(rebased_data));
+ return Status::OK();
+}
+
+// Compact sliced nested parents so recursive normalization sees only visible
child values.
+Status compact_lance_nested_array(const std::shared_ptr<arrow::Array>& array,
+ std::shared_ptr<arrow::Array>* compacted) {
+ switch (array->type_id()) {
+ case arrow::Type::LIST:
+ return compact_lance_offset_array<arrow::ListArray, int32_t>(array,
compacted);
+ case arrow::Type::LARGE_LIST:
+ return compact_lance_offset_array<arrow::LargeListArray,
int64_t>(array, compacted);
+ case arrow::Type::MAP:
+ return compact_lance_offset_array<arrow::MapArray, int32_t>(array,
compacted);
+ case arrow::Type::FIXED_SIZE_LIST: {
+ const auto list =
std::dynamic_pointer_cast<arrow::FixedSizeListArray>(array);
+ if (list == nullptr) {
+ return Status::InvalidArgument("invalid sliced Lance fixed-size
list array: {}",
+ array->type()->ToString());
+ }
+ const auto child_begin = list->value_offset(0);
+ const auto child_length = list->length() * list->value_length();
+ const auto& values = list->values();
+ if (child_begin < 0 || child_length < 0 || child_begin >
values->length() - child_length) {
+ return Status::InvalidArgument(
+ "invalid sliced Lance fixed-size list range [{}, {}) for
child length {}",
+ child_begin, child_begin + child_length, values->length());
+ }
+ if (array->offset() == 0 && child_begin == 0 && child_length ==
values->length()) {
+ *compacted = array;
+ return Status::OK();
+ }
+ std::shared_ptr<arrow::ArrayData> rebased_data;
+ RETURN_IF_ERROR(rebase_lance_parent_data(array, &rebased_data));
+ rebased_data->child_data[0] = values->Slice(child_begin,
child_length)->data();
+ *compacted = arrow::MakeArray(std::move(rebased_data));
+ return Status::OK();
+ }
+ case arrow::Type::STRUCT: {
+ const auto struct_array =
std::dynamic_pointer_cast<arrow::StructArray>(array);
+ if (struct_array == nullptr) {
+ return Status::InvalidArgument("invalid sliced Lance struct array:
{}",
+ array->type()->ToString());
+ }
+ bool requires_compaction = array->offset() != 0;
+ for (const auto& child : array->data()->child_data) {
+ requires_compaction |= child->length != array->length();
+ }
+ if (!requires_compaction) {
+ *compacted = array;
+ return Status::OK();
+ }
+ std::shared_ptr<arrow::ArrayData> rebased_data;
+ RETURN_IF_ERROR(rebase_lance_parent_data(array, &rebased_data));
+ for (int child_idx = 0; child_idx <
static_cast<int>(struct_array->fields().size());
+ ++child_idx) {
+ rebased_data->child_data[child_idx] =
struct_array->field(child_idx)->data();
+ }
+ *compacted = arrow::MakeArray(std::move(rebased_data));
+ return Status::OK();
+ }
+ default:
+ *compacted = array;
+ return Status::OK();
+ }
+}
+
} // namespace
+// Normalize nested BFloat16 arrays while preserving offsets and null bitmaps.
+Status normalize_lance_arrow_array(const std::shared_ptr<arrow::Field>& field,
+ const std::shared_ptr<arrow::Array>& array,
+ std::shared_ptr<arrow::Array>* normalized) {
+ DORIS_CHECK(field != nullptr);
+ DORIS_CHECK(array != nullptr);
+ DORIS_CHECK(normalized != nullptr);
+
+ LanceExtensionKind extension_kind;
+ std::shared_ptr<arrow::DataType> storage_type;
+ RETURN_IF_ERROR(get_lance_extension(field, &extension_kind,
&storage_type));
+
+ auto storage_array = array;
+ if (array->type_id() == arrow::Type::EXTENSION) {
+ const auto extension_array =
std::dynamic_pointer_cast<arrow::ExtensionArray>(array);
+ if (extension_array == nullptr) {
+ return Status::InvalidArgument("invalid Arrow extension array for
Lance field '{}'",
+ field->name());
+ }
+ storage_array = extension_array->storage();
+ }
+ if (storage_array->type_id() != storage_type->id()) {
+ return Status::InvalidArgument(
+ "Lance field '{}' storage type {} does not match array type
{}", field->name(),
+ storage_type->ToString(), storage_array->type()->ToString());
+ }
+ if (extension_kind == LanceExtensionKind::BFLOAT16) {
+ return convert_bfloat16_array(storage_array, normalized);
+ }
+
+ const auto& child_fields = storage_type->fields();
+ const auto& child_data = storage_array->data()->child_data;
+ if (child_fields.empty()) {
+ *normalized = std::move(storage_array);
Review Comment:
[P2] Rebase sliced scalar storage before handing it to Doris. Vector
pagination forwards a nonzero offset to Lance/DataFusion, whose limit stream
can return a sliced RecordBatch; for an already-present projected/refine
Duration or `arrow.json` column, this leaf branch returns the storage unchanged
while `_fill_block_from_record_batch()` always reads `[0,row_count)`. With
default Arrow validation, the Nullable/Number/JSON SerDes reject the nonzero
offset; without it, Duration and 32-bit JSON storage index buffers from row
zero and return preceding values. This is distinct from the existing
hidden-child BFloat16 allocation thread. Please add a sliced Duration/JSON
materialization test with leading sentinels and either compact here or make the
SerDes offset-aware.
--
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]