This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new 5445bdf5b72 [Feature](lance) Support additional Arrow and Lance data
types (#67325)
5445bdf5b72 is described below
commit 5445bdf5b721e631bf7a68d67060b0192998f6ec
Author: zhaobo wang <[email protected]>
AuthorDate: Mon Sep 14 09:46:11 2026 +0800
[Feature](lance) Support additional Arrow and Lance data types (#67325)
### What problem does this PR solve?
Issue Number: close #66496
Problem Summary:
The Lance reader previously reported several Arrow and Lance-specific
types as `UNSUPPORTED`, preventing Doris from correctly discovering
schemas or reading these columns.
This PR adds Doris-side support for:
- Arrow Null as Doris `NULL`
- Arrow Duration as Doris `BIGINT`
- Arrow and Lance JSON extensions as Doris `JSON`
- Lance BFloat16 as Doris `FLOAT`
Implementation details:
- FE recognizes the supported Arrow and Lance extension metadata and
validates each extension's physical storage type.
- BE maps the same logical types during schema discovery and scan
execution.
- BFloat16 values are widened to Float32 without precision loss,
including values nested in arrays and other complex types.
- Ordinary nested Arrow columns bypass reconstruction when no BFloat16
or registered extension array requires normalization.
- Only top-level Arrow Null fields are supported. Nested Null fields
remain unsupported because the current complex-type deserialization path
cannot handle them safely.
### Release note
Add Doris support for Arrow Null and Duration types, and Lance JSON and
BFloat16 extensions.
### Check List (For Author)
- Test
- [x] Regression test
- [x] Unit Test
- [ ] Manual test (add detailed scripts or steps below)
- [ ] No need to test or manual test
Coverage includes:
- FE type-conversion tests for Null, Duration, JSON, and BFloat16.
- BE schema and value tests, including nested BFloat16 and the
nested-column no-op normalization path.
- `test_lance_catalog_all_types`
- `test_lance_s3_tvf`
Validation:
- Production-source diff checks passed.
- A complete local FE test run is blocked because the standalone
`fe-core` build cannot resolve Doris internal SNAPSHOT dependencies.
- A complete local BE test run is blocked by the stale macOS build
directory and third-party headers; CI will run the full suite.
- Behavior changed:
- [ ] No.
- [x] Yes. Supported Lance schemas are mapped to Doris types instead of
`UNSUPPORTED`.
- Does this need documentation?
- [ ] No.
- [x] Yes. A follow-up documentation PR will be submitted separately.
### Check List (For Reviewer who merge this PR)
- [ ] Confirm the release note
- [ ] Confirm test cases
- [ ] Confirm document
- [ ] Add branch pick label
---------
Co-authored-by: wangzhaobo <[email protected]>
Co-authored-by: wangzhaobo957-cloud <wangzhaobo@bytedance>
---
be/src/format_v2/lance/lance_reader_helper.cpp | 535 +++++++++++++-
be/src/format_v2/lance/lance_reader_helper.h | 17 +
be/src/format_v2/table/lance_reader.cpp | 26 +-
be/src/service/internal_service.cpp | 6 +-
be/test/format_v2/table/lance_reader_test.cpp | 765 ++++++++++++++++++++-
.../doris/datasource/FederationBackendPolicy.java | 14 +-
.../doris/datasource/lance/LanceTypeConverter.java | 78 ++-
.../datasource/lance/source/LanceScanNode.java | 44 ++
.../doris/datasource/tvf/source/TVFScanNode.java | 28 +-
.../org/apache/doris/system/BeSelectionPolicy.java | 10 +
.../ExternalFileTableValuedFunction.java | 11 +
.../tablefunction/LocalTableValuedFunction.java | 34 +-
.../datasource/lance/LanceTypeConverterTest.java | 98 ++-
.../datasource/lance/source/LanceScanNodeTest.java | 21 +
.../datasource/tvf/source/TVFScanNodeTest.java | 66 ++
.../apache/doris/system/SystemInfoServiceTest.java | 24 +
.../ExternalFileTableValuedFunctionTest.java | 60 ++
.../lance/test_lance_catalog_all_types.out | 16 +-
.../external_table_p0/lance/test_lance_s3_tvf.out | 17 +-
.../lance/test_lance_catalog_all_types.groovy | 28 +-
.../lance/test_lance_s3_tvf.groovy | 16 +-
21 files changed, 1798 insertions(+), 116 deletions(-)
diff --git a/be/src/format_v2/lance/lance_reader_helper.cpp
b/be/src/format_v2/lance/lance_reader_helper.cpp
index 446a4b40179..c47e67a38a1 100644
--- a/be/src/format_v2/lance/lance_reader_helper.cpp
+++ b/be/src/format_v2/lance/lance_reader_helper.cpp
@@ -17,14 +17,20 @@
#include "format_v2/lance/lance_reader_helper.h"
+#include <arrow/array.h>
+#include <arrow/builder.h>
+#include <arrow/extension_type.h>
#include <arrow/type.h>
#include <arrow/util/key_value_metadata.h>
#include <fmt/format.h>
#include <lance/lance.h>
+#include <bit>
#include <limits>
+#include <optional>
#include <unordered_set>
+#include "common/config.h"
#include "common/logging.h"
#include "core/data_type/data_type_array.h"
#include "core/data_type/data_type_factory.hpp"
@@ -32,11 +38,22 @@
#include "core/data_type/data_type_nothing.h"
#include "core/data_type/data_type_nullable.h"
#include "core/data_type/data_type_struct.h"
+#include "core/data_type_serde/arrow_validation.h"
+#include "exec/common/endian.h"
namespace doris::format::lance {
namespace {
constexpr std::string_view ARROW_EXTENSION_NAME = "ARROW:extension:name";
+constexpr std::string_view ARROW_JSON_EXTENSION = "arrow.json";
+constexpr std::string_view LANCE_JSON_EXTENSION = "lance.json";
+constexpr std::string_view LANCE_BFLOAT16_EXTENSION = "lance.bfloat16";
+
+enum class LanceExtensionKind {
+ NONE,
+ JSON,
+ BFLOAT16,
+};
int arrow_time_precision(arrow::TimeUnit::type unit) {
switch (unit) {
@@ -51,26 +68,79 @@ int arrow_time_precision(arrow::TimeUnit::type unit) {
return 6;
}
-Status check_arrow_field_semantics(const std::shared_ptr<arrow::Field>& field)
{
+// Extract and validate extension names from metadata and registered Arrow
types.
+Status get_lance_extension(const std::shared_ptr<arrow::Field>& field,
+ LanceExtensionKind* extension_kind,
+ std::shared_ptr<arrow::DataType>* storage_type) {
+ DORIS_CHECK(field != nullptr);
+ DORIS_CHECK(extension_kind != nullptr);
+ DORIS_CHECK(storage_type != nullptr);
+ *extension_kind = LanceExtensionKind::NONE;
+ *storage_type = field->type();
+
+ std::optional<std::string> extension_name;
if (field->HasMetadata()) {
- const auto extension_name =
field->metadata()->Get(ARROW_EXTENSION_NAME);
- if (extension_name.ok() && !extension_name.ValueUnsafe().empty()) {
+ const auto metadata_name =
field->metadata()->Get(ARROW_EXTENSION_NAME);
+ if (metadata_name.ok() && !metadata_name.ValueUnsafe().empty()) {
+ extension_name = metadata_name.ValueUnsafe();
+ }
+ }
+ if (field->type()->id() == arrow::Type::EXTENSION) {
+ const auto extension_type =
std::dynamic_pointer_cast<arrow::ExtensionType>(field->type());
+ if (extension_type == nullptr) {
+ return Status::InvalidArgument("invalid Arrow extension type for
Lance field '{}'",
+ field->name());
+ }
+ if (extension_name.has_value() && *extension_name !=
extension_type->extension_name()) {
+ return Status::InvalidArgument(
+ "conflicting Arrow extension names for Lance field '{}':
'{}' and '{}'",
+ field->name(), *extension_name,
extension_type->extension_name());
+ }
+ extension_name = extension_type->extension_name();
+ *storage_type = extension_type->storage_type();
+ }
+ if (!extension_name.has_value()) {
+ return Status::OK();
+ }
+
+ if (*extension_name == ARROW_JSON_EXTENSION) {
+ if ((*storage_type)->id() != arrow::Type::STRING &&
+ (*storage_type)->id() != arrow::Type::LARGE_STRING) {
return Status::NotSupported(
- "unsupported Lance Arrow extension type '{}' for field
'{}'",
- extension_name.ValueUnsafe(), field->name());
+ "Arrow JSON extension for Lance field '{}' requires UTF8
storage, got {}",
+ field->name(), (*storage_type)->ToString());
}
+ *extension_kind = LanceExtensionKind::JSON;
+ return Status::OK();
}
- if (field->type()->id() == arrow::Type::DICTIONARY) {
- return Status::NotSupported("unsupported Lance Arrow dictionary type
for field '{}': {}",
- field->name(), field->type()->ToString());
+ if (*extension_name == LANCE_JSON_EXTENSION) {
+ if ((*storage_type)->id() != arrow::Type::LARGE_BINARY) {
+ return Status::NotSupported(
+ "Lance JSON extension for field '{}' requires LARGE_BINARY
storage, got {}",
+ field->name(), (*storage_type)->ToString());
+ }
+ *extension_kind = LanceExtensionKind::JSON;
+ return Status::OK();
}
- return Status::OK();
+ if (*extension_name == LANCE_BFLOAT16_EXTENSION) {
+ if ((*storage_type)->id() != arrow::Type::FIXED_SIZE_BINARY ||
+
std::static_pointer_cast<arrow::FixedSizeBinaryType>(*storage_type)->byte_width()
!=
+ 2) {
+ return Status::NotSupported(
+ "Lance BFloat16 extension for field '{}' requires
FIXED_SIZE_BINARY(2) "
+ "storage, got {}",
+ field->name(), (*storage_type)->ToString());
+ }
+ *extension_kind = LanceExtensionKind::BFLOAT16;
+ return Status::OK();
+ }
+ return Status::NotSupported("unsupported Lance Arrow extension type '{}'
for field '{}'",
+ *extension_name, field->name());
}
+// Map an Arrow field to a Doris type, allowing Doris NULL only at the top
level.
Status arrow_field_to_doris_type(const std::shared_ptr<arrow::Field>& field,
- DataTypePtr* doris_type) {
- RETURN_IF_ERROR(check_arrow_field_semantics(field));
- const auto& arrow_type = field->type();
+ DataTypePtr* doris_type, bool allow_null) {
const auto nullable_primitive = [&](PrimitiveType type, int precision = 0,
int scale = 0,
int len = -1) {
*doris_type =
@@ -78,7 +148,28 @@ Status arrow_field_to_doris_type(const
std::shared_ptr<arrow::Field>& field,
return Status::OK();
};
+ LanceExtensionKind extension_kind;
+ std::shared_ptr<arrow::DataType> arrow_type;
+ RETURN_IF_ERROR(get_lance_extension(field, &extension_kind, &arrow_type));
+ switch (extension_kind) {
+ case LanceExtensionKind::JSON:
+ return nullable_primitive(TYPE_JSONB);
+ case LanceExtensionKind::BFLOAT16:
+ return nullable_primitive(TYPE_FLOAT);
+ case LanceExtensionKind::NONE:
+ break;
+ }
+ if (arrow_type->id() == arrow::Type::DICTIONARY) {
+ return Status::NotSupported("unsupported Lance Arrow dictionary type
for field '{}': {}",
+ field->name(), arrow_type->ToString());
+ }
+
switch (arrow_type->id()) {
+ case arrow::Type::NA:
+ return allow_null ? nullable_primitive(TYPE_NULL)
+ : Status::NotSupported(
+ "nested Arrow null type is unsupported for
Lance field '{}'",
+ field->name());
case arrow::Type::BOOL:
return nullable_primitive(TYPE_BOOLEAN);
case arrow::Type::INT8:
@@ -122,6 +213,8 @@ Status arrow_field_to_doris_type(const
std::shared_ptr<arrow::Field>& field,
const auto doris_type = timestamp->timezone().empty() ?
TYPE_DATETIMEV2 : TYPE_TIMESTAMPTZ;
return nullable_primitive(doris_type, 0,
arrow_time_precision(timestamp->unit()));
}
+ case arrow::Type::DURATION:
+ return nullable_primitive(TYPE_BIGINT);
case arrow::Type::DECIMAL128:
case arrow::Type::DECIMAL256: {
const auto decimal =
std::static_pointer_cast<arrow::DecimalType>(arrow_type);
@@ -144,17 +237,16 @@ Status arrow_field_to_doris_type(const
std::shared_ptr<arrow::Field>& field,
case arrow::Type::FIXED_SIZE_LIST: {
const auto list =
std::static_pointer_cast<arrow::BaseListType>(arrow_type);
DataTypePtr value_type;
- RETURN_IF_ERROR(arrow_field_to_doris_type(list->value_field(),
&value_type));
+ RETURN_IF_ERROR(arrow_field_to_doris_type(list->value_field(),
&value_type, false));
*doris_type =
make_nullable(std::make_shared<DataTypeArray>(value_type));
return Status::OK();
}
case arrow::Type::MAP: {
const auto map = std::static_pointer_cast<arrow::MapType>(arrow_type);
- RETURN_IF_ERROR(check_arrow_field_semantics(map->value_field()));
DataTypePtr key_type;
DataTypePtr item_type;
- RETURN_IF_ERROR(arrow_field_to_doris_type(map->key_field(),
&key_type));
- RETURN_IF_ERROR(arrow_field_to_doris_type(map->item_field(),
&item_type));
+ RETURN_IF_ERROR(arrow_field_to_doris_type(map->key_field(), &key_type,
false));
+ RETURN_IF_ERROR(arrow_field_to_doris_type(map->item_field(),
&item_type, false));
*doris_type = make_nullable(std::make_shared<DataTypeMap>(key_type,
item_type));
return Status::OK();
}
@@ -166,7 +258,7 @@ Status arrow_field_to_doris_type(const
std::shared_ptr<arrow::Field>& field,
field_names.reserve(struct_type->num_fields());
for (const auto& child : struct_type->fields()) {
DataTypePtr field_type;
- RETURN_IF_ERROR(arrow_field_to_doris_type(child, &field_type));
+ RETURN_IF_ERROR(arrow_field_to_doris_type(child, &field_type,
false));
field_types.emplace_back(std::move(field_type));
field_names.emplace_back(child->name());
}
@@ -178,8 +270,415 @@ 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,
+ arrow::MemoryPool* memory_pool,
+ std::shared_ptr<arrow::Array>* normalized) {
+ DORIS_CHECK(array != nullptr);
+ DORIS_CHECK(memory_pool != 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(memory_pool);
+ 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();
+}
+
+// Rebuild a nested Arrow type after one or more child arrays changed physical
type.
+Status set_lance_nested_type(std::string_view field_name,
+ const std::shared_ptr<arrow::DataType>&
source_type,
+ const arrow::FieldVector& child_fields,
+ std::shared_ptr<arrow::ArrayData>* data) {
+ switch (source_type->id()) {
+ case arrow::Type::LIST:
+ (*data)->type = arrow::list(child_fields[0]);
+ break;
+ case arrow::Type::LARGE_LIST:
+ (*data)->type = arrow::large_list(child_fields[0]);
+ break;
+ case arrow::Type::FIXED_SIZE_LIST:
+ (*data)->type = arrow::fixed_size_list(
+ child_fields[0],
+
std::static_pointer_cast<arrow::FixedSizeListType>(source_type)->list_size());
+ break;
+ case arrow::Type::STRUCT:
+ (*data)->type = arrow::struct_(child_fields);
+ break;
+ case arrow::Type::MAP: {
+ const auto map_type =
std::static_pointer_cast<arrow::MapType>(source_type);
+ auto normalized_type = arrow::MapType::Make(child_fields[0],
map_type->keys_sorted());
+ if (!normalized_type.ok()) {
+ return Status::InvalidArgument("normalize Lance map field '{}'
failed: {}", field_name,
+ normalized_type.status().message());
+ }
+ (*data)->type = std::move(normalized_type).ValueUnsafe();
+ break;
+ }
+ default:
+ return Status::InvalidArgument("Lance field '{}' has unexpected
child-bearing type {}",
+ field_name, source_type->ToString());
+ }
+ return Status::OK();
+}
+
+// Check whether an Arrow type tree contains a registered extension wrapper.
+bool type_contains_registered_extension(const
std::shared_ptr<arrow::DataType>& type) {
+ if (type->id() == arrow::Type::EXTENSION) {
+ return true;
+ }
+ for (const auto& field : type->fields()) {
+ if (type_contains_registered_extension(field->type())) {
+ return true;
+ }
+ }
+ return false;
+}
+
+// Remove registered ExtensionArray wrappers only along extension-bearing
branches.
+Status unwrap_lance_extension_arrays(const std::shared_ptr<arrow::DataType>&
expected_type,
+ const std::shared_ptr<arrow::Array>&
array,
+ std::shared_ptr<arrow::Array>* unwrapped)
{
+ DORIS_CHECK(expected_type != nullptr);
+ DORIS_CHECK(array != nullptr);
+ DORIS_CHECK(unwrapped != nullptr);
+
+ auto storage_array = array;
+ auto expected_storage_type = expected_type;
+ if (expected_type->id() == arrow::Type::EXTENSION) {
+ const auto extension_type =
std::dynamic_pointer_cast<arrow::ExtensionType>(expected_type);
+ if (extension_type == nullptr) {
+ return Status::InvalidArgument("invalid expected Arrow extension
type {}",
+ expected_type->ToString());
+ }
+ expected_storage_type = extension_type->storage_type();
+ }
+ 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: {}",
+ array->type()->ToString());
+ }
+ storage_array = extension_array->storage();
+ }
+
+ const auto& child_data = storage_array->data()->child_data;
+ const auto& child_fields = expected_storage_type->fields();
+ if (child_data.empty()) {
+ *unwrapped = std::move(storage_array);
+ return Status::OK();
+ }
+ if (child_fields.size() != child_data.size()) {
+ return Status::InvalidArgument(
+ "Arrow array type {} has {} child fields but its data has {}
children",
+ storage_array->type()->ToString(), child_fields.size(),
child_data.size());
+ }
+
+ std::shared_ptr<arrow::ArrayData> unwrapped_data;
+ arrow::FieldVector unwrapped_fields;
+ for (size_t child_idx = 0; child_idx < child_data.size(); ++child_idx) {
+ if
(!type_contains_registered_extension(child_fields[child_idx]->type())) {
+ continue;
+ }
+ auto child_array = arrow::MakeArray(child_data[child_idx]);
+ std::shared_ptr<arrow::Array> unwrapped_child;
+
RETURN_IF_ERROR(unwrap_lance_extension_arrays(child_fields[child_idx]->type(),
child_array,
+ &unwrapped_child));
+ if (unwrapped_child.get() == child_array.get()) {
+ continue;
+ }
+ if (unwrapped_data == nullptr) {
+ unwrapped_data = storage_array->data()->Copy();
+ unwrapped_fields = storage_array->type()->fields();
+ }
+ unwrapped_data->child_data[child_idx] = unwrapped_child->data();
+ unwrapped_fields[child_idx] =
+ unwrapped_fields[child_idx]->WithType(unwrapped_child->type());
+ }
+ if (unwrapped_data == nullptr) {
+ *unwrapped = std::move(storage_array);
+ return Status::OK();
+ }
+ RETURN_IF_ERROR(
+ set_lance_nested_type("", storage_array->type(), unwrapped_fields,
&unwrapped_data));
+ *unwrapped = arrow::MakeArray(std::move(unwrapped_data));
+ return Status::OK();
+}
+
+// Materialize the visible range into an offset-zero Arrow array for Doris
SerDes.
+Status compact_lance_array(const std::shared_ptr<arrow::Array>& array,
+ arrow::MemoryPool* memory_pool,
+ std::shared_ptr<arrow::Array>* compacted) {
+ DORIS_CHECK(memory_pool != nullptr);
+ const auto validation = array->ValidateFull();
+ if (!validation.ok()) {
+ return Status::InvalidArgument("validate sliced Lance array failed:
{}",
+ validation.message());
+ }
+ auto builder_result = arrow::MakeBuilder(array->type(), memory_pool);
+ if (!builder_result.ok()) {
+ return Status::InternalError("create sliced Lance array builder
failed: {}",
+ builder_result.status().message());
+ }
+ auto builder = std::move(builder_result).ValueUnsafe();
+ auto arrow_status = builder->Reserve(array->length());
+ if (!arrow_status.ok()) {
+ return Status::InternalError("reserve sliced Lance array builder
failed: {}",
+ arrow_status.message());
+ }
+ arrow_status = builder->AppendArraySlice(*array->data(), 0,
array->length());
+ if (!arrow_status.ok()) {
+ return Status::InternalError("copy sliced Lance array failed: {}",
arrow_status.message());
+ }
+ arrow_status = builder->Finish(compacted);
+ if (!arrow_status.ok()) {
+ return Status::InternalError("finish sliced Lance array copy failed:
{}",
+ arrow_status.message());
+ }
+ if ((*compacted)->offset() != 0) {
+ return Status::InternalError("compacted Lance array retained offset
{}",
+ (*compacted)->offset());
+ }
+ return Status::OK();
+}
+
+// Compact a sliced variable-offset parent and all of its visible descendants.
+template <typename ArrayType>
+Status compact_lance_offset_array(const std::shared_ptr<arrow::Array>& array,
+ arrow::MemoryPool* memory_pool,
+ 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();
+ }
+ return compact_lance_array(array, memory_pool, compacted);
+}
+
+// Compact sliced arrays and nested children only when Doris cannot consume
their current layout.
+Status compact_lance_array_if_needed(const std::shared_ptr<arrow::Array>&
array,
+ arrow::MemoryPool* memory_pool,
+ std::shared_ptr<arrow::Array>* compacted)
{
+ switch (array->type_id()) {
+ case arrow::Type::LIST:
+ return compact_lance_offset_array<arrow::ListArray>(array,
memory_pool, compacted);
+ case arrow::Type::LARGE_LIST:
+ return compact_lance_offset_array<arrow::LargeListArray>(array,
memory_pool, compacted);
+ case arrow::Type::MAP:
+ return compact_lance_offset_array<arrow::MapArray>(array, memory_pool,
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();
+ }
+ return compact_lance_array(array, memory_pool, compacted);
+ }
+ 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();
+ }
+ return compact_lance_array(array, memory_pool, compacted);
+ }
+ default:
+ if (array->offset() == 0) {
+ *compacted = array;
+ return Status::OK();
+ }
+ return compact_lance_array(array, memory_pool, compacted);
+ }
+}
+
} // namespace
+// Normalize Lance extensions and materialize sliced arrays for Doris Arrow
SerDes.
+Status normalize_lance_arrow_array(const std::shared_ptr<arrow::Field>& field,
+ const std::shared_ptr<arrow::Array>& array,
+ arrow::MemoryPool* memory_pool,
+ std::shared_ptr<arrow::Array>* normalized) {
+ DORIS_CHECK(field != nullptr);
+ DORIS_CHECK(array != nullptr);
+ DORIS_CHECK(memory_pool != 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 (type_contains_registered_extension(field->type())) {
+ std::shared_ptr<arrow::Array> unwrapped_array;
+ RETURN_IF_ERROR(unwrap_lance_extension_arrays(field->type(), array,
&unwrapped_array));
+ storage_array = std::move(unwrapped_array);
+ }
+ 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, memory_pool, normalized);
+ }
+
+ std::shared_ptr<arrow::Array> compacted_array;
+ RETURN_IF_ERROR(compact_lance_array_if_needed(storage_array, memory_pool,
&compacted_array));
+ storage_array = std::move(compacted_array);
+
+ 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);
+ return Status::OK();
+ }
+ if (child_fields.size() != child_data.size()) {
+ return Status::InvalidArgument(
+ "Lance field '{}' has {} child fields but its Arrow array has
{} children",
+ field->name(), child_fields.size(), child_data.size());
+ }
+
+ bool requires_normalization = false;
+ for (const auto& child_field : child_fields) {
+ bool child_required = false;
+ RETURN_IF_ERROR(field_requires_lance_normalization(child_field,
&child_required));
+ if (child_required) {
+ requires_normalization = true;
+ break;
+ }
+ }
+ if (!requires_normalization) {
+ *normalized = std::move(storage_array);
+ return Status::OK();
+ }
+
+ arrow::FieldVector normalized_fields;
+ std::shared_ptr<arrow::ArrayData> normalized_data;
+ for (size_t child_idx = 0; child_idx < child_fields.size(); ++child_idx) {
+ bool child_required = false;
+ RETURN_IF_ERROR(
+ field_requires_lance_normalization(child_fields[child_idx],
&child_required));
+ if (!child_required) {
+ continue;
+ }
+ auto child_array =
arrow::MakeArray(storage_array->data()->child_data[child_idx]);
+ std::shared_ptr<arrow::Array> normalized_child;
+ RETURN_IF_ERROR(normalize_lance_arrow_array(child_fields[child_idx],
child_array,
+ memory_pool,
&normalized_child));
+ if (normalized_child.get() == child_array.get()) {
+ continue;
+ }
+ if (normalized_data == nullptr) {
+ normalized_data = storage_array->data()->Copy();
+ normalized_fields = storage_array->type()->fields();
+ }
+ normalized_data->child_data[child_idx] = normalized_child->data();
+ normalized_fields[child_idx] =
+
normalized_fields[child_idx]->WithType(normalized_child->type());
+ }
+ if (normalized_data == nullptr) {
+ *normalized = std::move(storage_array);
+ return Status::OK();
+ }
+
+ RETURN_IF_ERROR(set_lance_nested_type(field->name(),
storage_array->type(), normalized_fields,
+ &normalized_data));
+ *normalized = arrow::MakeArray(std::move(normalized_data));
+ return Status::OK();
+}
+
+#ifdef BE_TEST
+// Expose Lance Arrow normalization for allocation-sensitive unit tests.
+Status normalize_lance_arrow_array_for_test(const
std::shared_ptr<arrow::Field>& field,
+ const
std::shared_ptr<arrow::Array>& array,
+ std::shared_ptr<arrow::Array>*
normalized,
+ arrow::MemoryPool* memory_pool) {
+ return normalize_lance_arrow_array(
+ field, array, memory_pool != nullptr ? memory_pool :
arrow::default_memory_pool(),
+ normalized);
+}
+#endif
+
void LanceDatasetDeleter::operator()(LanceDataset* dataset) const {
lance_dataset_close(dataset);
}
@@ -267,7 +766,7 @@ Status convert_arrow_schema_to_doris(const
std::shared_ptr<arrow::Schema>& arrow
return Status::InvalidArgument("duplicate Lance schema column:
{}", field->name());
}
DataTypePtr doris_type;
- const auto type_status = arrow_field_to_doris_type(field, &doris_type);
+ const auto type_status = arrow_field_to_doris_type(field, &doris_type,
true);
if (type_status.is<ErrorCode::NOT_IMPLEMENTED_ERROR>()) {
parsed_types.emplace_back(std::make_shared<DataTypeNothing>());
} else {
diff --git a/be/src/format_v2/lance/lance_reader_helper.h
b/be/src/format_v2/lance/lance_reader_helper.h
index 689e4f4fd9f..8a8fb5a4f59 100644
--- a/be/src/format_v2/lance/lance_reader_helper.h
+++ b/be/src/format_v2/lance/lance_reader_helper.h
@@ -33,6 +33,9 @@ struct LanceDataset;
struct LanceScanner;
namespace arrow {
+class Array;
+class Field;
+class MemoryPool;
class Schema;
} // namespace arrow
@@ -63,6 +66,20 @@ Status parse_fragment_ids(const TLanceFileDesc&
lance_params, std::vector<uint64
Status parse_index_segment_uuids(const TLanceFileDesc& lance_params,
std::vector<uint8_t>* segment_uuids, size_t*
segment_count);
+// Normalize Lance extension arrays into Arrow arrays supported by Doris.
+Status normalize_lance_arrow_array(const std::shared_ptr<arrow::Field>& field,
+ const std::shared_ptr<arrow::Array>& array,
+ arrow::MemoryPool* memory_pool,
+ std::shared_ptr<arrow::Array>* normalized);
+
+#ifdef BE_TEST
+// Expose Lance Arrow normalization for allocation-sensitive unit tests.
+Status normalize_lance_arrow_array_for_test(const
std::shared_ptr<arrow::Field>& field,
+ const
std::shared_ptr<arrow::Array>& array,
+ std::shared_ptr<arrow::Array>*
normalized,
+ arrow::MemoryPool* memory_pool =
nullptr);
+#endif
+
// Convert every top-level field without discarding unsupported columns.
Malformed schemas still
// return an error and leave both output vectors unchanged. DataTypeNothing is
the local sentinel
// for a valid Arrow field whose logical type Doris does not support.
diff --git a/be/src/format_v2/table/lance_reader.cpp
b/be/src/format_v2/table/lance_reader.cpp
index a331dd231db..a9fc65cdc08 100644
--- a/be/src/format_v2/table/lance_reader.cpp
+++ b/be/src/format_v2/table/lance_reader.cpp
@@ -37,7 +37,9 @@
#include "exec/common/endian.h"
#include "format_v2/lance/lance_reader_helper.h"
#include "format_v2/lance/lance_runtime_filter_helper.h"
+#include "runtime/exec_env.h"
#include "runtime/file_scan_profile.h"
+#include "runtime/runtime_state.h"
#include "storage/utils.h"
namespace doris::format::lance {
@@ -63,6 +65,17 @@ Status import_dataset_schema(LanceDataset* dataset,
std::shared_ptr<arrow::Schem
return Status::OK();
}
+// Return the query-tracked Arrow pool, falling back only for standalone
unit-test states.
+arrow::MemoryPool* get_lance_arrow_memory_pool(RuntimeState* runtime_state) {
+ if (runtime_state != nullptr && runtime_state->exec_env() != nullptr) {
+ auto* memory_pool = runtime_state->exec_env()->arrow_memory_pool();
+ if (memory_pool != nullptr) {
+ return memory_pool;
+ }
+ }
+ return arrow::default_memory_pool();
+}
+
} // namespace
LanceTableReader::~LanceTableReader() {
@@ -1208,11 +1221,20 @@ Status LanceTableReader::_fill_block_from_record_batch(
}
const auto output_idx = output_it->second;
try {
+ const auto& arrow_column = record_batch->column(arrow_idx);
+ if (arrow_column->type_id() == arrow::Type::NA) {
+ columns[output_idx]->insert_many_defaults(row_count);
+ continue;
+ }
+ std::shared_ptr<arrow::Array> normalized_column;
+ RETURN_IF_ERROR(normalize_lance_arrow_array(field, arrow_column,
+
get_lance_arrow_memory_pool(_runtime_state),
+ &normalized_column));
RETURN_IF_ERROR(columns_guard.get_datatype_by_position(output_idx)
->get_serde()
->read_column_from_arrow(*columns[output_idx],
-
record_batch->column(arrow_idx).get(),
- 0, row_count,
_ctz));
+
normalized_column.get(), 0, row_count,
+ _ctz));
} catch (const Exception& e) {
return Status::InternalError("convert Lance Arrow column '{}'
failed: {}",
field->name(), e.what());
diff --git a/be/src/service/internal_service.cpp
b/be/src/service/internal_service.cpp
index 96142dae0b6..aa0b37e3c10 100644
--- a/be/src/service/internal_service.cpp
+++ b/be/src/service/internal_service.cpp
@@ -915,7 +915,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()) {
+ PTypeNode* node = type_desc->add_types();
+ node->set_type(TTypeNodeType::SCALAR);
+
node->mutable_scalar_type()->set_type(TPrimitiveType::NULL_TYPE);
+ } else if (col_type->get_primitive_type() == INVALID_TYPE) {
PTypeNode* node = type_desc->add_types();
node->set_type(TTypeNodeType::SCALAR);
node->mutable_scalar_type()->set_type(TPrimitiveType::UNSUPPORTED);
diff --git a/be/test/format_v2/table/lance_reader_test.cpp
b/be/test/format_v2/table/lance_reader_test.cpp
index 89a244a61a7..cfcb4597963 100644
--- a/be/test/format_v2/table/lance_reader_test.cpp
+++ b/be/test/format_v2/table/lance_reader_test.cpp
@@ -21,7 +21,10 @@
#include <arrow/array/builder_decimal.h>
#include <arrow/array/builder_primitive.h>
#include <arrow/array/util.h>
+#include <arrow/builder.h>
#include <arrow/c/bridge.h>
+#include <arrow/extension/json.h>
+#include <arrow/memory_pool.h>
#include <arrow/record_batch.h>
#include <arrow/type.h>
#include <arrow/util/decimal.h>
@@ -1630,6 +1633,49 @@ DataTypePtr nullable_type(PrimitiveType type, int
precision = 0, int scale = 0)
return DataTypeFactory::instance().create_data_type(type, true, precision,
scale);
}
+// Creates an Arrow Duration array with one value and one null.
+std::shared_ptr<arrow::Array> make_duration_array(arrow::TimeUnit::type unit,
int64_t value) {
+ arrow::DurationBuilder builder(arrow::duration(unit),
arrow::default_memory_pool());
+ EXPECT_TRUE(builder.Append(value).ok());
+ EXPECT_TRUE(builder.AppendNull().ok());
+ std::shared_ptr<arrow::DurationArray> array;
+ EXPECT_TRUE(builder.Finish(&array).ok());
+ return array;
+}
+
+// Creates a Lance BFloat16 array with one vector and one null.
+std::shared_ptr<arrow::Array> make_bfloat16_vector_array() {
+ auto values =
std::make_shared<arrow::FixedSizeBinaryBuilder>(arrow::fixed_size_binary(2));
+ arrow::FixedSizeListBuilder builder(arrow::default_memory_pool(), values,
2);
+ EXPECT_TRUE(builder.Append().ok());
+ const std::array<uint8_t, 2> one {0x80, 0x3F};
+ const std::array<uint8_t, 2> two {0x00, 0x40};
+ EXPECT_TRUE(values->Append(one.data()).ok());
+ EXPECT_TRUE(values->Append(two.data()).ok());
+ EXPECT_TRUE(builder.AppendNull().ok());
+ std::shared_ptr<arrow::FixedSizeListArray> array;
+ EXPECT_TRUE(builder.Finish(&array).ok());
+ return array;
+}
+
+// Creates a scalar Lance BFloat16 array from raw 16-bit values.
+std::shared_ptr<arrow::FixedSizeBinaryArray> make_bfloat16_array(
+ const std::vector<std::optional<uint16_t>>& values) {
+ arrow::FixedSizeBinaryBuilder builder(arrow::fixed_size_binary(2));
+ for (const auto& value : values) {
+ if (!value.has_value()) {
+ EXPECT_TRUE(builder.AppendNull().ok());
+ continue;
+ }
+ std::array<uint8_t, sizeof(uint16_t)> bytes {};
+ LittleEndian::Store16(bytes.data(), *value);
+ EXPECT_TRUE(builder.Append(bytes.data()).ok());
+ }
+ std::shared_ptr<arrow::FixedSizeBinaryArray> array;
+ EXPECT_TRUE(builder.Finish(&array).ok());
+ return array;
+}
+
std::pair<size_t, size_t> array_range(const ColumnArray& array, size_t row) {
const auto& offsets = array.get_offsets();
return {row == 0 ? 0 : static_cast<size_t>(offsets[row - 1]),
@@ -1680,17 +1726,25 @@ TEST(LanceTableReaderSchemaTest,
FetchesSchemaWithoutFragmentIdsOrScanInitializa
assert_cast<const DataTypeVarbinary&>(*binary_type).len());
}
-TEST(LanceTableReaderSchemaTest,
PreservesUnsupportedFieldsAndExtensionSemantics) {
- const auto extension_metadata =
+// Verifies the additional mappings and preserves unknown extensions as
unsupported.
+TEST(LanceTableReaderSchemaTest,
MapsAdditionalTypesAndPreservesUnknownExtensions) {
+ const auto unknown_extension_metadata =
arrow::KeyValueMetadata::Make({"ARROW:extension:name"},
{"doris.test.extension"});
- const auto extension_item =
- arrow::field("item",
arrow::float32())->WithMetadata(extension_metadata);
+ const auto json_extension_metadata =
+ arrow::KeyValueMetadata::Make({"ARROW:extension:name"},
{"arrow.json"});
+ const auto bfloat16_extension_metadata =
+ arrow::KeyValueMetadata::Make({"ARROW:extension:name"},
{"lance.bfloat16"});
+ const auto bfloat16_item = arrow::field("item",
arrow::fixed_size_binary(2))
+
->WithMetadata(bfloat16_extension_metadata);
const auto arrow_schema = arrow::schema({
arrow::field("row_id", arrow::int64()),
+ arrow::field("null_value", arrow::null()),
arrow::field("duration", arrow::duration(arrow::TimeUnit::MILLI)),
- arrow::field("json",
arrow::utf8())->WithMetadata(extension_metadata),
+ arrow::field("json",
arrow::utf8())->WithMetadata(json_extension_metadata),
+ arrow::field("bfloat16_vector",
arrow::fixed_size_list(bfloat16_item, 4)),
arrow::field("dictionary", arrow::dictionary(arrow::int16(),
arrow::utf8())),
- arrow::field("nested_extension", arrow::list(extension_item)),
+ arrow::field("unknown_extension", arrow::utf8())
+ ->WithMetadata(unknown_extension_metadata),
arrow::field("name", arrow::utf8()),
});
@@ -1698,19 +1752,704 @@ TEST(LanceTableReaderSchemaTest,
PreservesUnsupportedFieldsAndExtensionSemantics
std::vector<DataTypePtr> column_types;
ASSERT_TRUE(convert_arrow_schema_to_doris(arrow_schema, &column_names,
&column_types).ok());
- EXPECT_EQ((std::vector<std::string> {"row_id", "duration", "json",
"dictionary",
- "nested_extension", "name"}),
+ EXPECT_EQ((std::vector<std::string> {"row_id", "null_value", "duration",
"json",
+ "bfloat16_vector", "dictionary",
"unknown_extension",
+ "name"}),
column_names);
ASSERT_EQ(column_names.size(), column_types.size());
for (const auto& column_type : column_types) {
ASSERT_NE(nullptr, column_type);
}
EXPECT_EQ(TYPE_BIGINT, column_types[0]->get_primitive_type());
- EXPECT_EQ(INVALID_TYPE, column_types[1]->get_primitive_type());
- EXPECT_EQ(INVALID_TYPE, column_types[2]->get_primitive_type());
- EXPECT_EQ(INVALID_TYPE, column_types[3]->get_primitive_type());
- EXPECT_EQ(INVALID_TYPE, column_types[4]->get_primitive_type());
- EXPECT_EQ(TYPE_STRING, column_types[5]->get_primitive_type());
+ EXPECT_TRUE(column_types[1]->is_null_literal());
+ EXPECT_EQ(TYPE_BIGINT, column_types[2]->get_primitive_type());
+ EXPECT_EQ(TYPE_JSONB, column_types[3]->get_primitive_type());
+ ASSERT_EQ(TYPE_ARRAY, column_types[4]->get_primitive_type());
+ const auto& bfloat16_array =
+ assert_cast<const
DataTypeArray&>(*remove_nullable(column_types[4]));
+ EXPECT_EQ(TYPE_FLOAT,
bfloat16_array.get_nested_type()->get_primitive_type());
+ EXPECT_EQ(INVALID_TYPE, column_types[5]->get_primitive_type());
+ EXPECT_EQ(INVALID_TYPE, column_types[6]->get_primitive_type());
+ EXPECT_EQ(TYPE_STRING, column_types[7]->get_primitive_type());
+}
+
+// Verifies malformed storage for known extensions remains unsupported.
+TEST(LanceTableReaderSchemaTest, RejectsMalformedKnownExtensionStorage) {
+ const auto json_extension_metadata =
+ arrow::KeyValueMetadata::Make({"ARROW:extension:name"},
{"arrow.json"});
+ const auto bfloat16_extension_metadata =
+ arrow::KeyValueMetadata::Make({"ARROW:extension:name"},
{"lance.bfloat16"});
+ const auto arrow_schema = arrow::schema({
+ arrow::field("json",
arrow::binary())->WithMetadata(json_extension_metadata),
+ arrow::field("bfloat16", arrow::fixed_size_binary(4))
+ ->WithMetadata(bfloat16_extension_metadata),
+ });
+
+ std::vector<std::string> column_names;
+ std::vector<DataTypePtr> column_types;
+ ASSERT_TRUE(convert_arrow_schema_to_doris(arrow_schema, &column_names,
&column_types).ok());
+ ASSERT_EQ(2, column_types.size());
+ for (const auto& column_type : column_types) {
+ ASSERT_NE(nullptr, column_type);
+ EXPECT_EQ(INVALID_TYPE, column_type->get_primitive_type());
+ }
+}
+
+// Verifies nested Null fields remain unsupported.
+TEST(LanceTableReaderSchemaTest, MarksNestedNullTypesAsUnsupported) {
+ const auto arrow_schema = arrow::schema({
+ arrow::field("null_list", arrow::list(arrow::field("item",
arrow::null()))),
+ arrow::field("null_struct", arrow::struct_({arrow::field("value",
arrow::null())})),
+ });
+
+ std::vector<std::string> column_names;
+ std::vector<DataTypePtr> column_types;
+ ASSERT_TRUE(convert_arrow_schema_to_doris(arrow_schema, &column_names,
&column_types).ok());
+
+ EXPECT_EQ((std::vector<std::string> {"null_list", "null_struct"}),
column_names);
+ ASSERT_EQ(2, column_types.size());
+ for (const auto& column_type : column_types) {
+ ASSERT_NE(nullptr, column_type);
+ EXPECT_EQ(INVALID_TYPE, column_type->get_primitive_type());
+ }
+}
+
+// Verifies values, nullability, and precision when reading the additional
types.
+TEST(LanceTableReaderTypeTest, ReadsAdditionalArrowAndLanceTypes) {
+ const auto json_extension_metadata =
+ arrow::KeyValueMetadata::Make({"ARROW:extension:name"},
{"arrow.json"});
+ const auto bfloat16_extension_metadata =
+ arrow::KeyValueMetadata::Make({"ARROW:extension:name"},
{"lance.bfloat16"});
+ const auto bfloat16_item = arrow::field("item",
arrow::fixed_size_binary(2))
+
->WithMetadata(bfloat16_extension_metadata);
+ const auto schema = arrow::schema({
+ arrow::field("null_value", arrow::null()),
+ arrow::field("duration_s",
arrow::duration(arrow::TimeUnit::SECOND)),
+ arrow::field("duration_ms",
arrow::duration(arrow::TimeUnit::MILLI)),
+ arrow::field("duration_us",
arrow::duration(arrow::TimeUnit::MICRO)),
+ arrow::field("duration_ns",
arrow::duration(arrow::TimeUnit::NANO)),
+ arrow::field("json_value",
arrow::utf8())->WithMetadata(json_extension_metadata),
+ arrow::field("bfloat16_vector",
arrow::fixed_size_list(bfloat16_item, 2)),
+ });
+
+ arrow::StringBuilder json_builder;
+ ASSERT_TRUE(json_builder.Append(R"({"engine":"doris"})").ok());
+ ASSERT_TRUE(json_builder.AppendNull().ok());
+ std::shared_ptr<arrow::StringArray> json_array;
+ ASSERT_TRUE(json_builder.Finish(&json_array).ok());
+
+ const auto record_batch =
+ arrow::RecordBatch::Make(schema, 2,
+ {
+
std::make_shared<arrow::NullArray>(2),
+
make_duration_array(arrow::TimeUnit::SECOND, 1),
+
make_duration_array(arrow::TimeUnit::MILLI, 1000),
+
make_duration_array(arrow::TimeUnit::MICRO, 1000000),
+
make_duration_array(arrow::TimeUnit::NANO, 1000000000),
+ json_array,
+ make_bfloat16_vector_array(),
+ });
+
+ const auto bfloat16_array_type =
+
make_nullable(std::make_shared<DataTypeArray>(nullable_type(TYPE_FLOAT)));
+ const Columns columns {
+ projected_column("null_value", nullable_type(TYPE_NULL)),
+ projected_column("duration_s", TYPE_BIGINT, true),
+ projected_column("duration_ms", TYPE_BIGINT, true),
+ projected_column("duration_us", TYPE_BIGINT, true),
+ projected_column("duration_ns", TYPE_BIGINT, true),
+ projected_column("json_value", TYPE_JSONB, true),
+ projected_column("bfloat16_vector", bfloat16_array_type),
+ };
+ TQueryGlobals query_globals;
+ RuntimeState state(query_globals);
+ RuntimeProfile profile("lance_additional_types");
+ TFileScanRangeParams scan_params;
+ LanceTableReader reader;
+ ASSERT_TRUE(init_reader(&reader, columns, &state, &profile,
&scan_params).ok());
+
+ Block block;
+ add_output_columns(&block, columns);
+ size_t rows = 0;
+ ASSERT_TRUE(reader._fill_block_from_record_batch(record_batch, &block,
&rows).ok());
+ ASSERT_EQ(2, rows);
+
+ const auto& null_values = assert_cast<const
ColumnNullable&>(*block.get_by_position(0).column);
+ EXPECT_EQ((ColumnUInt8::Container {1, 1}),
null_values.get_null_map_data());
+
+ const std::array<int64_t, 4> expected_durations {1, 1000, 1000000,
1000000000};
+ for (size_t column_idx = 0; column_idx < expected_durations.size();
++column_idx) {
+ const auto& duration =
+ assert_cast<const
ColumnNullable&>(*block.get_by_position(column_idx + 1).column);
+ const auto& values = assert_cast<const
ColumnInt64&>(duration.get_nested_column());
+ EXPECT_EQ(0, duration.get_null_map_data()[0]);
+ EXPECT_EQ(1, duration.get_null_map_data()[1]);
+ EXPECT_EQ(expected_durations[column_idx], values.get_data()[0]);
+ }
+
+ const auto& json_values = assert_cast<const
ColumnNullable&>(*block.get_by_position(5).column);
+ EXPECT_EQ(0, json_values.get_null_map_data()[0]);
+ EXPECT_EQ(1, json_values.get_null_map_data()[1]);
+ EXPECT_EQ(R"({"engine":"doris"})", columns[5].type->to_string(json_values,
0));
+
+ const auto& vectors = assert_cast<const
ColumnNullable&>(*block.get_by_position(6).column);
+ EXPECT_EQ(0, vectors.get_null_map_data()[0]);
+ EXPECT_EQ(1, vectors.get_null_map_data()[1]);
+ const auto& vector_values = assert_cast<const
ColumnArray&>(vectors.get_nested_column());
+ EXPECT_EQ((ColumnArray::Offsets64 {2, 4}), vector_values.get_offsets());
+ const auto& bfloat16_values = assert_cast<const
ColumnNullable&>(vector_values.get_data());
+ const auto& floats = assert_cast<const
ColumnFloat32&>(bfloat16_values.get_nested_column());
+ EXPECT_FLOAT_EQ(1.0F, floats.get_data()[0]);
+ EXPECT_FLOAT_EQ(2.0F, floats.get_data()[1]);
+ EXPECT_TRUE(reader.close().ok());
+}
+
+// Verifies sliced scalar arrays are compacted before Doris reads rows from
offset zero.
+TEST(LanceTableReaderTypeTest, ReadsSlicedDurationAndJsonValues) {
+ const auto json_extension_metadata =
+ arrow::KeyValueMetadata::Make({"ARROW:extension:name"},
{"arrow.json"});
+ const auto duration_type = arrow::duration(arrow::TimeUnit::MILLI);
+ const auto schema = arrow::schema({
+ arrow::field("duration", duration_type),
+ arrow::field("json_value",
arrow::utf8())->WithMetadata(json_extension_metadata),
+ });
+
+ arrow::DurationBuilder duration_builder(duration_type,
arrow::default_memory_pool());
+ ASSERT_TRUE(duration_builder.Append(111).ok());
+ ASSERT_TRUE(duration_builder.Append(222).ok());
+ ASSERT_TRUE(duration_builder.AppendNull().ok());
+ ASSERT_TRUE(duration_builder.Append(444).ok());
+ std::shared_ptr<arrow::DurationArray> durations;
+ ASSERT_TRUE(duration_builder.Finish(&durations).ok());
+
+ arrow::StringBuilder json_builder;
+ ASSERT_TRUE(json_builder.Append(R"({"sentinel":true})").ok());
+ ASSERT_TRUE(json_builder.Append(R"({"row":1})").ok());
+ ASSERT_TRUE(json_builder.Append(R"({"row":2})").ok());
+ ASSERT_TRUE(json_builder.Append(R"({"tail":true})").ok());
+ std::shared_ptr<arrow::StringArray> json_values;
+ ASSERT_TRUE(json_builder.Finish(&json_values).ok());
+
+ const auto record_batch =
+ arrow::RecordBatch::Make(schema, 2, {durations->Slice(1, 2),
json_values->Slice(1, 2)});
+ const Columns columns {
+ projected_column("duration", TYPE_BIGINT, true),
+ projected_column("json_value", TYPE_JSONB, true),
+ };
+ TQueryGlobals query_globals;
+ RuntimeState state(query_globals);
+ RuntimeProfile profile("lance_sliced_scalars");
+ TFileScanRangeParams scan_params;
+ LanceTableReader reader;
+ ASSERT_TRUE(init_reader(&reader, columns, &state, &profile,
&scan_params).ok());
+
+ Block block;
+ add_output_columns(&block, columns);
+ size_t rows = 0;
+ ASSERT_TRUE(reader._fill_block_from_record_batch(record_batch, &block,
&rows).ok());
+ ASSERT_EQ(2, rows);
+
+ const auto& duration = assert_cast<const
ColumnNullable&>(*block.get_by_position(0).column);
+ const auto& duration_values = assert_cast<const
ColumnInt64&>(duration.get_nested_column());
+ EXPECT_EQ((ColumnUInt8::Container {0, 1}), duration.get_null_map_data());
+ EXPECT_EQ(222, duration_values.get_data()[0]);
+
+ const auto& json = assert_cast<const
ColumnNullable&>(*block.get_by_position(1).column);
+ EXPECT_EQ((ColumnUInt8::Container {0, 0}), json.get_null_map_data());
+ EXPECT_EQ(R"({"row":1})", columns[1].type->to_string(json, 0));
+ EXPECT_EQ(R"({"row":2})", columns[1].type->to_string(json, 1));
+ EXPECT_TRUE(reader.close().ok());
+}
+
+// Verifies sliced nested registered JSON arrays are unwrapped before parent
compaction.
+TEST(LanceTableReaderTypeTest, ReadsRegisteredJsonNestedInSlicedList) {
+ const auto json_type = arrow::extension::json();
+ const auto item_field = arrow::field("item", json_type);
+ const auto list_type = arrow::list(item_field);
+
+ arrow::StringBuilder json_builder;
+ ASSERT_TRUE(json_builder.Append(R"({"sentinel":true})").ok());
+ ASSERT_TRUE(json_builder.Append(R"({"row":1})").ok());
+ ASSERT_TRUE(json_builder.Append(R"({"row":2})").ok());
+ ASSERT_TRUE(json_builder.Append(R"({"tail":true})").ok());
+ std::shared_ptr<arrow::StringArray> json_storage;
+ ASSERT_TRUE(json_builder.Finish(&json_storage).ok());
+ const auto json_values = arrow::ExtensionType::WrapArray(json_type,
json_storage);
+
+ arrow::Int32Builder offsets_builder;
+ ASSERT_TRUE(offsets_builder.AppendValues({0, 1, 3, 4}).ok());
+ std::shared_ptr<arrow::Int32Array> offsets;
+ ASSERT_TRUE(offsets_builder.Finish(&offsets).ok());
+ auto list_result = arrow::ListArray::FromArrays(list_type, *offsets,
*json_values);
+ ASSERT_TRUE(list_result.ok()) << list_result.status().ToString();
+ const auto input = std::move(list_result).ValueUnsafe()->Slice(1, 1);
+
+ std::shared_ptr<arrow::Array> normalized;
+ const auto status =
normalize_lance_arrow_array_for_test(arrow::field("values", list_type),
+ input,
&normalized);
+ ASSERT_TRUE(status.ok()) << status.to_string();
+ const auto normalized_list =
std::dynamic_pointer_cast<arrow::ListArray>(normalized);
+ ASSERT_NE(nullptr, normalized_list);
+ EXPECT_EQ(0, normalized_list->offset());
+ ASSERT_EQ(2, normalized_list->values()->length());
+ const auto normalized_json =
+
std::dynamic_pointer_cast<arrow::StringArray>(normalized_list->values());
+ ASSERT_NE(nullptr, normalized_json);
+ EXPECT_EQ(R"({"row":1})", normalized_json->GetString(0));
+ EXPECT_EQ(R"({"row":2})", normalized_json->GetString(1));
+
+ const auto schema = arrow::schema({arrow::field("values", list_type)});
+ const auto record_batch = arrow::RecordBatch::Make(schema, 1, {input});
+ const auto doris_list_type =
+
make_nullable(std::make_shared<DataTypeArray>(nullable_type(TYPE_JSONB)));
+ const Columns columns {projected_column("values", doris_list_type)};
+ TQueryGlobals query_globals;
+ RuntimeState state(query_globals);
+ RuntimeProfile profile("lance_sliced_nested_json");
+ TFileScanRangeParams scan_params;
+ LanceTableReader reader;
+ ASSERT_TRUE(init_reader(&reader, columns, &state, &profile,
&scan_params).ok());
+
+ Block block;
+ add_output_columns(&block, columns);
+ size_t rows = 0;
+ ASSERT_TRUE(reader._fill_block_from_record_batch(record_batch, &block,
&rows).ok());
+ ASSERT_EQ(1, rows);
+
+ const auto& nullable_list =
+ assert_cast<const
ColumnNullable&>(*block.get_by_position(0).column);
+ const auto& list = assert_cast<const
ColumnArray&>(nullable_list.get_nested_column());
+ EXPECT_EQ((ColumnArray::Offsets64 {2}), list.get_offsets());
+ const auto& items = assert_cast<const ColumnNullable&>(list.get_data());
+ const auto json_item_type = nullable_type(TYPE_JSONB);
+ EXPECT_EQ(R"({"row":1})", json_item_type->to_string(items, 0));
+ EXPECT_EQ(R"({"row":2})", json_item_type->to_string(items, 1));
+ EXPECT_TRUE(reader.close().ok());
+}
+
+// Verifies BFloat16 special values retain their exact Float32 bit patterns.
+TEST(LanceTableReaderTypeTest, ConvertsBFloat16SpecialValuesExactly) {
+ const auto metadata =
+ arrow::KeyValueMetadata::Make({"ARROW:extension:name"},
{"lance.bfloat16"});
+ const auto field = arrow::field("value",
arrow::fixed_size_binary(2))->WithMetadata(metadata);
+ const std::vector<std::optional<uint16_t>> input_bits {0x0000, 0x8000,
0x3F80, 0xBF80,
+ 0x7F80, 0xFF80,
0x7FC1, std::nullopt};
+ const std::array<uint32_t, 7> expected_bits {0x00000000, 0x80000000,
0x3F800000, 0xBF800000,
+ 0x7F800000, 0xFF800000,
0x7FC10000};
+
+ std::shared_ptr<arrow::Array> normalized;
+ const auto status = normalize_lance_arrow_array_for_test(field,
make_bfloat16_array(input_bits),
+ &normalized);
+ ASSERT_TRUE(status.ok()) << status.to_string();
+ const auto floats =
std::dynamic_pointer_cast<arrow::FloatArray>(normalized);
+ ASSERT_NE(nullptr, floats);
+ ASSERT_EQ(input_bits.size(), static_cast<size_t>(floats->length()));
+ for (size_t index = 0; index < expected_bits.size(); ++index) {
+ ASSERT_FALSE(floats->IsNull(index));
+ EXPECT_EQ(expected_bits[index],
std::bit_cast<uint32_t>(floats->Value(index)));
+ }
+ EXPECT_TRUE(floats->IsNull(expected_bits.size()));
+}
+
+// Verifies normalization honors the supplied pool when widening and
compacting arrays.
+TEST(LanceTableReaderTypeTest, HonorsNormalizationMemoryPoolLimit) {
+ arrow::ProxyMemoryPool proxy_pool(arrow::default_memory_pool());
+ arrow::CappedMemoryPool capped_pool(&proxy_pool, 0);
+ std::shared_ptr<arrow::Array> normalized;
+
+ const auto bfloat16_metadata =
+ arrow::KeyValueMetadata::Make({"ARROW:extension:name"},
{"lance.bfloat16"});
+ const auto bfloat16_field =
+ arrow::field("value",
arrow::fixed_size_binary(2))->WithMetadata(bfloat16_metadata);
+ const auto bfloat16_status = normalize_lance_arrow_array_for_test(
+ bfloat16_field, make_bfloat16_array({0x3F80}), &normalized,
&capped_pool);
+ EXPECT_FALSE(bfloat16_status.ok());
+ EXPECT_NE(std::string::npos,
+ bfloat16_status.to_string().find("reserve Lance BFloat16 output
failed"));
+
+ const auto duration_type = arrow::duration(arrow::TimeUnit::MILLI);
+ arrow::DurationBuilder duration_builder(duration_type,
arrow::default_memory_pool());
+ ASSERT_TRUE(duration_builder.AppendValues({100, 200}).ok());
+ std::shared_ptr<arrow::DurationArray> durations;
+ ASSERT_TRUE(duration_builder.Finish(&durations).ok());
+ const auto duration_status =
+ normalize_lance_arrow_array_for_test(arrow::field("duration",
duration_type),
+ durations->Slice(1, 1),
&normalized, &capped_pool);
+ EXPECT_FALSE(duration_status.ok());
+ EXPECT_NE(std::string::npos,
+ duration_status.to_string().find("reserve sliced Lance array
builder failed"));
+}
+
+// Verifies nested BFloat16 fields are converted without rebuilding unaffected
sibling data.
+TEST(LanceTableReaderTypeTest, ConvertsBFloat16NestedInStruct) {
+ const auto metadata =
+ arrow::KeyValueMetadata::Make({"ARROW:extension:name"},
{"lance.bfloat16"});
+ const auto value_field =
+ arrow::field("value",
arrow::fixed_size_binary(2))->WithMetadata(metadata);
+ const auto id_field = arrow::field("id", arrow::int32());
+ const auto bfloat16_values = make_bfloat16_array(
+ {std::optional<uint16_t> {0x3F80}, std::optional<uint16_t>
{0xC020}});
+ arrow::Int32Builder id_builder;
+ ASSERT_TRUE(id_builder.AppendValues({7, 8}).ok());
+ std::shared_ptr<arrow::Int32Array> ids;
+ ASSERT_TRUE(id_builder.Finish(&ids).ok());
+ auto struct_result = arrow::StructArray::Make({bfloat16_values, ids},
{value_field, id_field});
+ ASSERT_TRUE(struct_result.ok()) << struct_result.status().ToString();
+ const auto input = std::move(struct_result).ValueUnsafe();
+ const auto field = arrow::field("record", input->type());
+
+ std::shared_ptr<arrow::Array> normalized;
+ const auto status = normalize_lance_arrow_array_for_test(field, input,
&normalized);
+ ASSERT_TRUE(status.ok()) << status.to_string();
+ ASSERT_NE(input.get(), normalized.get());
+ const auto normalized_struct =
std::dynamic_pointer_cast<arrow::StructArray>(normalized);
+ ASSERT_NE(nullptr, normalized_struct);
+ const auto floats =
std::dynamic_pointer_cast<arrow::FloatArray>(normalized_struct->field(0));
+ ASSERT_NE(nullptr, floats);
+ EXPECT_FLOAT_EQ(1.0F, floats->Value(0));
+ EXPECT_FLOAT_EQ(-2.5F, floats->Value(1));
+ EXPECT_EQ(ids->data().get(), normalized_struct->field(1)->data().get());
+}
+
+// Verifies sliced List normalization converts only values visible through
parent offsets.
+TEST(LanceTableReaderTypeTest, NormalizesVisibleBFloat16ValuesInSlicedList) {
+ const auto metadata =
+ arrow::KeyValueMetadata::Make({"ARROW:extension:name"},
{"lance.bfloat16"});
+ const auto item_field =
+ arrow::field("item",
arrow::fixed_size_binary(2))->WithMetadata(metadata);
+ const auto values = make_bfloat16_array({0x3F80, 0x4000, 0x4040, 0x4080,
0x40A0, 0x40C0});
+ arrow::Int32Builder offsets_builder;
+ ASSERT_TRUE(offsets_builder.AppendValues({0, 2, 5, 6}).ok());
+ std::shared_ptr<arrow::Int32Array> offsets;
+ ASSERT_TRUE(offsets_builder.Finish(&offsets).ok());
+ auto list_result = arrow::ListArray::FromArrays(arrow::list(item_field),
*offsets, *values);
+ ASSERT_TRUE(list_result.ok()) << list_result.status().ToString();
+ const auto input = std::move(list_result).ValueUnsafe()->Slice(1, 1);
+
+ std::shared_ptr<arrow::Array> normalized;
+ const auto status = normalize_lance_arrow_array_for_test(
+ arrow::field("values", arrow::list(item_field)), input,
&normalized);
+ ASSERT_TRUE(status.ok()) << status.to_string();
+ const auto list = std::dynamic_pointer_cast<arrow::ListArray>(normalized);
+ ASSERT_NE(nullptr, list);
+ EXPECT_EQ(0, list->offset());
+ ASSERT_EQ(1, list->length());
+ EXPECT_EQ(0, list->value_offset(0));
+ EXPECT_EQ(3, list->value_offset(1));
+ const auto floats =
std::dynamic_pointer_cast<arrow::FloatArray>(list->values());
+ ASSERT_NE(nullptr, floats);
+ ASSERT_EQ(3, floats->length());
+ EXPECT_FLOAT_EQ(3.0F, floats->Value(0));
+ EXPECT_FLOAT_EQ(4.0F, floats->Value(1));
+ EXPECT_FLOAT_EQ(5.0F, floats->Value(2));
+}
+
+// Verifies sliced LargeList normalization rebases 64-bit offsets and trims
child values.
+TEST(LanceTableReaderTypeTest,
NormalizesVisibleBFloat16ValuesInSlicedLargeList) {
+ const auto metadata =
+ arrow::KeyValueMetadata::Make({"ARROW:extension:name"},
{"lance.bfloat16"});
+ const auto item_field =
+ arrow::field("item",
arrow::fixed_size_binary(2))->WithMetadata(metadata);
+ const auto values = make_bfloat16_array({0x3F80, 0x4000, 0x4040, 0x4080,
0x40A0, 0x40C0});
+ arrow::Int64Builder offsets_builder;
+ ASSERT_TRUE(offsets_builder.AppendValues({0, 1, 3, 6}).ok());
+ std::shared_ptr<arrow::Int64Array> offsets;
+ ASSERT_TRUE(offsets_builder.Finish(&offsets).ok());
+ auto list_result =
+ arrow::LargeListArray::FromArrays(arrow::large_list(item_field),
*offsets, *values);
+ ASSERT_TRUE(list_result.ok()) << list_result.status().ToString();
+ const auto input = std::move(list_result).ValueUnsafe()->Slice(1, 1);
+
+ std::shared_ptr<arrow::Array> normalized;
+ const auto status = normalize_lance_arrow_array_for_test(
+ arrow::field("values", arrow::large_list(item_field)), input,
&normalized);
+ ASSERT_TRUE(status.ok()) << status.to_string();
+ const auto list =
std::dynamic_pointer_cast<arrow::LargeListArray>(normalized);
+ ASSERT_NE(nullptr, list);
+ EXPECT_EQ(0, list->offset());
+ ASSERT_EQ(1, list->length());
+ EXPECT_EQ(0, list->value_offset(0));
+ EXPECT_EQ(2, list->value_offset(1));
+ const auto floats =
std::dynamic_pointer_cast<arrow::FloatArray>(list->values());
+ ASSERT_NE(nullptr, floats);
+ ASSERT_EQ(2, floats->length());
+ EXPECT_FLOAT_EQ(2.0F, floats->Value(0));
+ EXPECT_FLOAT_EQ(3.0F, floats->Value(1));
+}
+
+// Verifies sliced FixedSizeList normalization trims values using the fixed
child width.
+TEST(LanceTableReaderTypeTest,
NormalizesVisibleBFloat16ValuesInSlicedFixedSizeList) {
+ const auto metadata =
+ arrow::KeyValueMetadata::Make({"ARROW:extension:name"},
{"lance.bfloat16"});
+ const auto item_field =
+ arrow::field("item",
arrow::fixed_size_binary(2))->WithMetadata(metadata);
+ const auto values =
+ make_bfloat16_array({0x3F80, 0x4000, 0x4040, 0x4080, 0x40A0,
0x40C0, 0x40E0, 0x4100});
+ auto list_result = arrow::FixedSizeListArray::FromArrays(values, 2);
+ ASSERT_TRUE(list_result.ok()) << list_result.status().ToString();
+ const auto input = std::move(list_result).ValueUnsafe()->Slice(2, 1);
+
+ std::shared_ptr<arrow::Array> normalized;
+ const auto status = normalize_lance_arrow_array_for_test(
+ arrow::field("values", arrow::fixed_size_list(item_field, 2)),
input, &normalized);
+ ASSERT_TRUE(status.ok()) << status.to_string();
+ const auto list =
std::dynamic_pointer_cast<arrow::FixedSizeListArray>(normalized);
+ ASSERT_NE(nullptr, list);
+ EXPECT_EQ(0, list->offset());
+ ASSERT_EQ(1, list->length());
+ const auto floats =
std::dynamic_pointer_cast<arrow::FloatArray>(list->values());
+ ASSERT_NE(nullptr, floats);
+ ASSERT_EQ(2, floats->length());
+ EXPECT_FLOAT_EQ(5.0F, floats->Value(0));
+ EXPECT_FLOAT_EQ(6.0F, floats->Value(1));
+}
+
+// Verifies sliced Map normalization trims entries while preserving visible
keys and values.
+TEST(LanceTableReaderTypeTest, NormalizesVisibleBFloat16ValuesInSlicedMap) {
+ const auto metadata =
+ arrow::KeyValueMetadata::Make({"ARROW:extension:name"},
{"lance.bfloat16"});
+ const auto item_field =
+ arrow::field("value",
arrow::fixed_size_binary(2))->WithMetadata(metadata);
+ const auto values = make_bfloat16_array({0x3F80, 0x4000, 0x4040, 0x4080,
0x40A0, 0x40C0});
+ arrow::Int32Builder offsets_builder;
+ ASSERT_TRUE(offsets_builder.AppendValues({0, 2, 5, 6}).ok());
+ std::shared_ptr<arrow::Int32Array> offsets;
+ ASSERT_TRUE(offsets_builder.Finish(&offsets).ok());
+ arrow::StringBuilder keys_builder;
+ ASSERT_TRUE(keys_builder.AppendValues({"k0", "k1", "k2", "k3", "k4",
"k5"}).ok());
+ std::shared_ptr<arrow::StringArray> keys;
+ ASSERT_TRUE(keys_builder.Finish(&keys).ok());
+ const auto map_type = arrow::map(arrow::utf8(), item_field);
+ auto map_result = arrow::MapArray::FromArrays(map_type, offsets, keys,
values);
+ ASSERT_TRUE(map_result.ok()) << map_result.status().ToString();
+ const auto input = std::move(map_result).ValueUnsafe()->Slice(1, 1);
+
+ std::shared_ptr<arrow::Array> normalized;
+ const auto status =
normalize_lance_arrow_array_for_test(arrow::field("values", map_type),
+ input,
&normalized);
+ ASSERT_TRUE(status.ok()) << status.to_string();
+ const auto map = std::dynamic_pointer_cast<arrow::MapArray>(normalized);
+ ASSERT_NE(nullptr, map);
+ EXPECT_EQ(0, map->offset());
+ ASSERT_EQ(1, map->length());
+ EXPECT_EQ(0, map->value_offset(0));
+ EXPECT_EQ(3, map->value_offset(1));
+ const auto normalized_keys =
std::dynamic_pointer_cast<arrow::StringArray>(map->keys());
+ ASSERT_NE(nullptr, normalized_keys);
+ ASSERT_EQ(3, normalized_keys->length());
+ EXPECT_EQ("k2", normalized_keys->GetString(0));
+ EXPECT_EQ("k3", normalized_keys->GetString(1));
+ EXPECT_EQ("k4", normalized_keys->GetString(2));
+ const auto floats =
std::dynamic_pointer_cast<arrow::FloatArray>(map->items());
+ ASSERT_NE(nullptr, floats);
+ ASSERT_EQ(3, floats->length());
+ EXPECT_FLOAT_EQ(3.0F, floats->Value(0));
+ EXPECT_FLOAT_EQ(4.0F, floats->Value(1));
+ EXPECT_FLOAT_EQ(5.0F, floats->Value(2));
+
+ const auto schema = arrow::schema({arrow::field("values", map_type)});
+ const auto record_batch = arrow::RecordBatch::Make(schema, 1, {input});
+ const auto doris_map_type = make_nullable(
+ std::make_shared<DataTypeMap>(nullable_type(TYPE_STRING),
nullable_type(TYPE_FLOAT)));
+ const Columns columns {projected_column("values", doris_map_type)};
+ TQueryGlobals query_globals;
+ RuntimeState state(query_globals);
+ RuntimeProfile profile("lance_sliced_map");
+ TFileScanRangeParams scan_params;
+ LanceTableReader reader;
+ ASSERT_TRUE(init_reader(&reader, columns, &state, &profile,
&scan_params).ok());
+
+ Block block;
+ add_output_columns(&block, columns);
+ size_t rows = 0;
+ ASSERT_TRUE(reader._fill_block_from_record_batch(record_batch, &block,
&rows).ok());
+ ASSERT_EQ(1, rows);
+
+ const auto& nullable_map = assert_cast<const
ColumnNullable&>(*block.get_by_position(0).column);
+ const auto& materialized_map = assert_cast<const
ColumnMap&>(nullable_map.get_nested_column());
+ EXPECT_EQ((ColumnArray::Offsets64 {3}), materialized_map.get_offsets());
+ const auto& materialized_keys = assert_cast<const
ColumnNullable&>(materialized_map.get_keys());
+ const auto& key_values =
+ assert_cast<const
ColumnString&>(materialized_keys.get_nested_column());
+ const auto& materialized_items =
+ assert_cast<const ColumnNullable&>(materialized_map.get_values());
+ const auto& item_values =
+ assert_cast<const
ColumnFloat32&>(materialized_items.get_nested_column());
+ ASSERT_EQ(3, key_values.size());
+ ASSERT_EQ(3, item_values.size());
+ EXPECT_EQ("k2", key_values.get_data_at(0).to_string());
+ EXPECT_EQ("k3", key_values.get_data_at(1).to_string());
+ EXPECT_EQ("k4", key_values.get_data_at(2).to_string());
+ EXPECT_FLOAT_EQ(3.0F, item_values.get_data()[0]);
+ EXPECT_FLOAT_EQ(4.0F, item_values.get_data()[1]);
+ EXPECT_FLOAT_EQ(5.0F, item_values.get_data()[2]);
+ EXPECT_TRUE(reader.close().ok());
+}
+
+// Verifies the Lance JSON extension reads LargeBinary values and nulls
through JSON SerDe.
+TEST(LanceTableReaderTypeTest, ReadsLanceJsonLargeBinaryValues) {
+ const auto metadata =
arrow::KeyValueMetadata::Make({"ARROW:extension:name"}, {"lance.json"});
+ const auto schema = arrow::schema(
+ {arrow::field("json_value",
arrow::large_binary())->WithMetadata(metadata)});
+ arrow::LargeBinaryBuilder builder;
+ ASSERT_TRUE(builder.Append(R"({"format":"lance","value":42})").ok());
+ ASSERT_TRUE(builder.AppendNull().ok());
+ std::shared_ptr<arrow::LargeBinaryArray> values;
+ ASSERT_TRUE(builder.Finish(&values).ok());
+ const auto record_batch = arrow::RecordBatch::Make(schema, 2, {values});
+
+ const Columns columns {projected_column("json_value", TYPE_JSONB, true)};
+ TQueryGlobals query_globals;
+ RuntimeState state(query_globals);
+ RuntimeProfile profile("lance_json_large_binary");
+ TFileScanRangeParams scan_params;
+ LanceTableReader reader;
+ ASSERT_TRUE(init_reader(&reader, columns, &state, &profile,
&scan_params).ok());
+
+ Block block;
+ add_output_columns(&block, columns);
+ size_t rows = 0;
+ ASSERT_TRUE(reader._fill_block_from_record_batch(record_batch, &block,
&rows).ok());
+ ASSERT_EQ(2, rows);
+ const auto& json_values = assert_cast<const
ColumnNullable&>(*block.get_by_position(0).column);
+ EXPECT_EQ((ColumnUInt8::Container {0, 1}),
json_values.get_null_map_data());
+ EXPECT_EQ(R"({"format":"lance","value":42})",
columns[0].type->to_string(json_values, 0));
+ EXPECT_TRUE(reader.close().ok());
+}
+
+// Verifies Duration values preserve signed 64-bit boundaries and nullability.
+TEST(LanceTableReaderTypeTest, ReadsDurationBoundaryValues) {
+ const auto duration_type = arrow::duration(arrow::TimeUnit::NANO);
+ const auto schema = arrow::schema({arrow::field("duration",
duration_type)});
+ arrow::DurationBuilder builder(duration_type,
arrow::default_memory_pool());
+ const std::array<int64_t, 4> expected
{std::numeric_limits<int64_t>::min(), -1, 0,
+
std::numeric_limits<int64_t>::max()};
+ ASSERT_TRUE(builder.AppendValues(expected.data(), expected.size()).ok());
+ ASSERT_TRUE(builder.AppendNull().ok());
+ std::shared_ptr<arrow::DurationArray> values;
+ ASSERT_TRUE(builder.Finish(&values).ok());
+ const auto record_batch = arrow::RecordBatch::Make(schema, 5, {values});
+
+ const Columns columns {projected_column("duration", TYPE_BIGINT, true)};
+ TQueryGlobals query_globals;
+ RuntimeState state(query_globals);
+ RuntimeProfile profile("lance_duration_boundaries");
+ TFileScanRangeParams scan_params;
+ LanceTableReader reader;
+ ASSERT_TRUE(init_reader(&reader, columns, &state, &profile,
&scan_params).ok());
+
+ Block block;
+ add_output_columns(&block, columns);
+ size_t rows = 0;
+ ASSERT_TRUE(reader._fill_block_from_record_batch(record_batch, &block,
&rows).ok());
+ ASSERT_EQ(5, rows);
+ const auto& durations = assert_cast<const
ColumnNullable&>(*block.get_by_position(0).column);
+ const auto& duration_values = assert_cast<const
ColumnInt64&>(durations.get_nested_column());
+ EXPECT_EQ((ColumnUInt8::Container {0, 0, 0, 0, 1}),
durations.get_null_map_data());
+ for (size_t index = 0; index < expected.size(); ++index) {
+ EXPECT_EQ(expected[index], duration_values.get_data()[index]);
+ }
+ EXPECT_TRUE(reader.close().ok());
+}
+
+// Verifies ordinary nested columns bypass Arrow reconstruction when no
BFloat16 exists.
+TEST(LanceTableReaderTypeTest,
KeepsOrdinaryNestedArrayUnchangedDuringNormalization) {
+ const auto nested_type = arrow::struct_({
+ arrow::field("items",
+ arrow::list(arrow::field(
+ "item", arrow::struct_({arrow::field("value",
arrow::int64())})))),
+ arrow::field("metadata", arrow::struct_({arrow::field("name",
arrow::utf8()),
+ arrow::field("enabled",
arrow::boolean())})),
+ });
+ const auto field = arrow::field("ordinary_nested", nested_type);
+ auto array_result = arrow::MakeArrayOfNull(nested_type, 3);
+ ASSERT_TRUE(array_result.ok()) << array_result.status().ToString();
+ const auto array = std::move(array_result).ValueUnsafe();
+
+ std::shared_ptr<arrow::Array> normalized;
+ const auto status = normalize_lance_arrow_array_for_test(field, array,
&normalized);
+ ASSERT_TRUE(status.ok()) << status.to_string();
+ EXPECT_EQ(array.get(), normalized.get());
+}
+
+// Verifies the additional types through the full scan path.
+TEST(LanceTableReaderTypeTest, ReadsAdditionalTypesFromCompatibilityFixture) {
+ const std::filesystem::path dataset_uri =
+
"./docker/thirdparties/docker-compose/iceberg/scripts/preinstalled_data/lance/"
+ "all_types.lance";
+ const auto bfloat16_array_type =
+
make_nullable(std::make_shared<DataTypeArray>(nullable_type(TYPE_FLOAT)));
+ const Columns columns {
+ projected_column("row_id", TYPE_BIGINT, false),
+ projected_column("null_col", nullable_type(TYPE_NULL)),
+ projected_column("duration_s_col", TYPE_BIGINT, true),
+ projected_column("duration_ms_col", TYPE_BIGINT, true),
+ projected_column("duration_us_col", TYPE_BIGINT, true),
+ projected_column("duration_ns_col", TYPE_BIGINT, true),
+ projected_column("json_col", TYPE_JSONB, true),
+ projected_column("bfloat16_vector_col", bfloat16_array_type),
+ };
+ TQueryOptions query_options;
+ query_options.__set_batch_size(4);
+ TQueryGlobals query_globals;
+ RuntimeState state(query_globals);
+ state.set_query_options(query_options);
+ RuntimeProfile profile("lance_additional_types_fixture");
+ TFileScanRangeParams scan_params;
+ LanceTableReader reader;
+ ASSERT_TRUE(init_reader(&reader, columns, &state, &profile,
&scan_params).ok());
+ ASSERT_TRUE(prepare_range(&reader,
make_latest_lance_range(dataset_uri)).ok());
+
+ Block block;
+ add_output_columns(&block, columns);
+ bool found = false;
+ bool eos = false;
+ while (!eos) {
+ ASSERT_TRUE(reader.get_block(&block, &eos).ok());
+ if (eos) {
+ continue;
+ }
+ const auto& row_ids = assert_cast<const
ColumnInt64&>(*block.get_by_position(0).column);
+ for (size_t row = 0; row < block.rows(); ++row) {
+ const auto& null_values =
+ assert_cast<const
ColumnNullable&>(*block.get_by_position(1).column);
+ EXPECT_EQ(1, null_values.get_null_map_data()[row]);
+ if (row_ids.get_data()[row] != 1) {
+ continue;
+ }
+ found = true;
+ const std::array<int64_t, 4> expected_durations {1, 1000, 1000000,
1000000000};
+ for (size_t duration_idx = 0; duration_idx <
expected_durations.size();
+ ++duration_idx) {
+ const auto& duration = assert_cast<const ColumnNullable&>(
+ *block.get_by_position(duration_idx + 2).column);
+ const auto& values = assert_cast<const
ColumnInt64&>(duration.get_nested_column());
+ EXPECT_EQ(0, duration.get_null_map_data()[row]);
+ EXPECT_EQ(expected_durations[duration_idx],
values.get_data()[row]);
+ }
+
+ const auto& json_values =
+ assert_cast<const
ColumnNullable&>(*block.get_by_position(6).column);
+ EXPECT_EQ(R"({"engine":"doris","format":"lance"})",
+ columns[6].type->to_string(json_values, row));
+
+ const auto& vectors =
+ assert_cast<const
ColumnNullable&>(*block.get_by_position(7).column);
+ const auto& vector_values =
+ assert_cast<const
ColumnArray&>(vectors.get_nested_column());
+ const auto [begin, end] = array_range(vector_values, row);
+ ASSERT_EQ(4, end - begin);
+ const auto& nullable_values =
+ assert_cast<const
ColumnNullable&>(vector_values.get_data());
+ const auto& floats =
+ assert_cast<const
ColumnFloat32&>(nullable_values.get_nested_column());
+ for (size_t index = 0; index < 4; ++index) {
+ EXPECT_FLOAT_EQ(static_cast<float>(index + 1),
floats.get_data()[begin + index]);
+ }
+ }
+ }
+ EXPECT_TRUE(found);
+ EXPECT_TRUE(reader.close().ok());
}
TEST(LanceTableReaderSchemaTest, FetchesNegativeScaleDecimalAsUnsupported) {
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/datasource/FederationBackendPolicy.java
b/fe/fe-core/src/main/java/org/apache/doris/datasource/FederationBackendPolicy.java
index 6fd43dfb26e..2ac01065bbb 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/datasource/FederationBackendPolicy.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/datasource/FederationBackendPolicy.java
@@ -156,13 +156,25 @@ public class FederationBackendPolicy {
}
public void init(List<String> preLocations) throws UserException {
+ init(preLocations, Collections.emptyList());
+ }
+
+ /** Initialize the policy with exactly one eligible backend. */
+ public void initWithBackendId(long backendId) throws UserException {
+ init(Collections.emptyList(), Collections.singletonList(backendId));
+ }
+
+ /** Build the standard external-scan policy with optional location and
backend-ID constraints. */
+ private void init(List<String> preLocations, List<Long> requiredBackendIds)
+ throws UserException {
// scan node is used for query
BeSelectionPolicy.Builder builder = new BeSelectionPolicy.Builder();
builder.needQueryAvailable()
.needLoadAvailable()
.preferComputeNode(Config.prefer_compute_node_for_external_table)
.assignExpectBeNum(Config.min_backend_num_for_external_table)
- .addPreLocations(preLocations);
+ .addPreLocations(preLocations)
+ .addRequiredBackendIds(requiredBackendIds);
init(builder.build());
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceTypeConverter.java
b/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceTypeConverter.java
index 2686fb9d576..77f23f04144 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceTypeConverter.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceTypeConverter.java
@@ -37,28 +37,59 @@ import java.util.List;
public final class LanceTypeConverter {
private static final int MAX_DECIMAL_PRECISION = 76;
private static final String ARROW_EXTENSION_NAME = "ARROW:extension:name";
+ private static final String ARROW_JSON_EXTENSION = "arrow.json";
+ private static final String LANCE_JSON_EXTENSION = "lance.json";
+ private static final String LANCE_BFLOAT16_EXTENSION = "lance.bfloat16";
private LanceTypeConverter() {
}
+ /** Converts Arrow fields exposed by Lance to Doris types. */
public static Type toDorisType(Field field) {
- // Arrow Java exposes unknown extension types through their storage
type and field
- // metadata. Treating the storage type as the logical type would make
DESC report
- // Blob, JSON, or BFloat16 as supported even though the scanner cannot
decode their
- // extension semantics. Dictionary arrays are likewise not decoded by
the BE reader.
+ return toDorisType(field, true);
+ }
+
+ /** Returns whether this field needs the current BE Lance materialization
logic. */
+ public static boolean requiresCurrentBeReader(Field field) {
+ ArrowType.ArrowTypeID typeId = field.getType().getTypeID();
+ if (typeId == ArrowType.ArrowTypeID.Null
+ || typeId == ArrowType.ArrowTypeID.Duration) {
+ return true;
+ }
+ String extensionName = field.getMetadata() == null
+ ? null : field.getMetadata().get(ARROW_EXTENSION_NAME);
+ if (ARROW_JSON_EXTENSION.equals(extensionName)
+ || LANCE_JSON_EXTENSION.equals(extensionName)
+ || LANCE_BFLOAT16_EXTENSION.equals(extensionName)) {
+ return true;
+ }
+ for (Field child : field.getChildren()) {
+ if (requiresCurrentBeReader(child)) {
+ return true;
+ }
+ }
+ return false;
+ }
+
+ /** Converts an Arrow field, allowing Doris NULL only at the top level. */
+ private static Type toDorisType(Field field, boolean allowNull) {
// TODO(lance): Dataset.getSchema() currently erases the Dictionary
marker, while
// Dataset.getLanceSchema() fails to convert a schema containing
Dictionary in the
// Lance 9.1.0-beta.3 Java SDK. Reject physical Dictionary columns
after that SDK
// conversion is fixed; an unmarked Int16 field cannot be
distinguished safely here.
String extensionName = field.getMetadata() == null
? null : field.getMetadata().get(ARROW_EXTENSION_NAME);
- if (field.getDictionary() != null
- || (extensionName != null && !extensionName.isEmpty())) {
+ if (field.getDictionary() != null) {
return Type.UNSUPPORTED;
}
+ if (extensionName != null && !extensionName.isEmpty()) {
+ return extensionType(field, extensionName);
+ }
ArrowType arrowType = field.getType();
switch (arrowType.getTypeID()) {
+ case Null:
+ return allowNull ? Type.NULL : Type.UNSUPPORTED;
case Bool:
return Type.BOOLEAN;
case Int:
@@ -91,6 +122,8 @@ public final class LanceTypeConverter {
return timeType((ArrowType.Time) arrowType);
case Timestamp:
return timestampType((ArrowType.Timestamp) arrowType);
+ case Duration:
+ return Type.BIGINT;
case Decimal:
ArrowType.Decimal decimal = (ArrowType.Decimal) arrowType;
if (decimal.getPrecision() <= 0 || decimal.getPrecision() >
MAX_DECIMAL_PRECISION
@@ -102,7 +135,7 @@ public final class LanceTypeConverter {
case LargeList:
case FixedSizeList:
requireChildren(field, 1);
- Type itemType = toDorisType(field.getChildren().get(0));
+ Type itemType = toDorisType(field.getChildren().get(0), false);
return itemType.isSupported() ? new ArrayType(itemType) :
Type.UNSUPPORTED;
case Map:
requireChildren(field, 1);
@@ -110,15 +143,15 @@ public final class LanceTypeConverter {
requireChildren(entries, 2);
Field key = entries.getChildren().get(0);
Field value = entries.getChildren().get(1);
- Type keyType = toDorisType(key);
- Type valueType = toDorisType(value);
+ Type keyType = toDorisType(key, false);
+ Type valueType = toDorisType(value, false);
return keyType.isSupported() && valueType.isSupported()
? new MapType(keyType, valueType, key.isNullable(),
value.isNullable())
: Type.UNSUPPORTED;
case Struct:
List<StructField> fields = new ArrayList<>();
for (Field child : field.getChildren()) {
- Type childType = toDorisType(child);
+ Type childType = toDorisType(child, false);
if (!childType.isSupported()) {
return Type.UNSUPPORTED;
}
@@ -133,6 +166,27 @@ public final class LanceTypeConverter {
}
}
+ /** Maps a known extension and validates its storage type. */
+ private static Type extensionType(Field field, String extensionName) {
+ ArrowType storageType = field.getType();
+ switch (extensionName) {
+ case ARROW_JSON_EXTENSION:
+ return storageType.getTypeID() == ArrowType.ArrowTypeID.Utf8
+ || storageType.getTypeID() ==
ArrowType.ArrowTypeID.LargeUtf8
+ ? Type.JSONB : Type.UNSUPPORTED;
+ case LANCE_JSON_EXTENSION:
+ return storageType.getTypeID() ==
ArrowType.ArrowTypeID.LargeBinary
+ ? Type.JSONB : Type.UNSUPPORTED;
+ case LANCE_BFLOAT16_EXTENSION:
+ return storageType.getTypeID() ==
ArrowType.ArrowTypeID.FixedSizeBinary
+ && ((ArrowType.FixedSizeBinary)
storageType).getByteWidth() == 2
+ ? Type.FLOAT : Type.UNSUPPORTED;
+ default:
+ return Type.UNSUPPORTED;
+ }
+ }
+
+ /** Maps an Arrow integer to the narrowest lossless Doris type. */
private static Type integerType(ArrowType.Int type) {
if (type.getIsSigned()) {
switch (type.getBitWidth()) {
@@ -162,6 +216,7 @@ public final class LanceTypeConverter {
}
}
+ /** Maps an Arrow timestamp while preserving precision and timezone
semantics. */
private static Type timestampType(ArrowType.Timestamp type) {
TimeUnit unit = type.getUnit();
int scale;
@@ -185,6 +240,7 @@ public final class LanceTypeConverter {
: ScalarType.createTimeStampTzType(scale);
}
+ /** Maps an Arrow time value to Doris TIMEV2. */
private static Type timeType(ArrowType.Time type) {
TimeUnit unit = type.getUnit();
switch (unit) {
@@ -203,6 +259,7 @@ public final class LanceTypeConverter {
}
}
+ /** Maps an Arrow date value to Doris DATEV2. */
private static Type dateType(ArrowType.Date type) {
DateUnit unit = type.getUnit();
switch (unit) {
@@ -214,6 +271,7 @@ public final class LanceTypeConverter {
}
}
+ /** Checks the child count of a nested Arrow field. */
private static void requireChildren(Field field, int expected) {
if (field.getChildren().size() != expected) {
throw new IllegalArgumentException("Invalid Arrow children for
Lance field '" + field.getName()
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LanceScanNode.java
b/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LanceScanNode.java
index 697ae298e9e..4bbde106e81 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LanceScanNode.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/datasource/lance/source/LanceScanNode.java
@@ -17,6 +17,7 @@
package org.apache.doris.datasource.lance.source;
+import org.apache.doris.analysis.SlotDescriptor;
import org.apache.doris.analysis.TupleDescriptor;
import org.apache.doris.catalog.Column;
import org.apache.doris.catalog.TableIf;
@@ -29,12 +30,14 @@ import org.apache.doris.datasource.lance.LanceExternalTable;
import org.apache.doris.datasource.lance.LanceFragmentInfo;
import org.apache.doris.datasource.lance.LanceIndexSegmentInfo;
import org.apache.doris.datasource.lance.LanceTableMetadata;
+import org.apache.doris.datasource.lance.LanceTypeConverter;
import org.apache.doris.datasource.mvcc.MvccSnapshot;
import org.apache.doris.planner.PlanNodeId;
import org.apache.doris.planner.ScanContext;
import org.apache.doris.qe.SessionVariable;
import org.apache.doris.spi.Split;
import org.apache.doris.statistics.StatisticalType;
+import org.apache.doris.system.Backend;
import org.apache.doris.thrift.TExplainLevel;
import org.apache.doris.thrift.TExternalSearchRequest;
import org.apache.doris.thrift.TFileFormatType;
@@ -49,13 +52,19 @@ import org.apache.doris.thrift.TVectorMetric;
import org.apache.doris.thrift.TVectorSearchOptions;
import org.apache.doris.thrift.TVectorSearchParams;
+import com.google.common.annotations.VisibleForTesting;
+import org.apache.arrow.vector.types.pojo.Field;
+
import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.Collections;
+import java.util.HashSet;
import java.util.LinkedHashMap;
import java.util.List;
+import java.util.Locale;
import java.util.Map;
import java.util.Optional;
+import java.util.Set;
import java.util.UUID;
/**
@@ -145,6 +154,8 @@ public class LanceScanNode extends FileQueryScanNode {
sourceColumns = lanceTable.getFullSchema(relationSnapshot);
}
super.doInitialize();
+ checkAdditionalTypeBackendCompatibility(
+ projectsCurrentReaderType(), backendPolicy.getBackends());
ExternalUtil.initSchemaInfo(params, -1L, sourceColumns);
if (searchKind != SearchKind.NORMAL) {
@@ -155,6 +166,39 @@ public class LanceScanNode extends FileQueryScanNode {
}
}
+ /** Checks whether any projected Lance column requires the current BE
reader. */
+ private boolean projectsCurrentReaderType() {
+ Set<String> projectedColumns = new HashSet<>();
+ for (SlotDescriptor slot : desc.getSlots()) {
+ if (slot.getColumn() != null) {
+
projectedColumns.add(slot.getColumn().getName().toLowerCase(Locale.ROOT));
+ }
+ }
+ for (Field field : plannedMetadata.getSchema().getFields()) {
+ if
(projectedColumns.contains(field.getName().toLowerCase(Locale.ROOT))
+ && LanceTypeConverter.requiresCurrentBeReader(field)) {
+ return true;
+ }
+ }
+ return false;
+ }
+
+ /** Rejects old smooth-upgrade source BEs for additional Lance type
projections. */
+ @VisibleForTesting
+ public static void checkAdditionalTypeBackendCompatibility(
+ boolean requiresCurrentReader, Iterable<Backend> backends) throws
UserException {
+ if (!requiresCurrentReader) {
+ return;
+ }
+ for (Backend backend : backends) {
+ if (backend.isSmoothUpgradeSrc()) {
+ throw new UserException(
+ "Additional Lance types are unavailable while backend "
+ + backend.getId() + " is a smooth upgrade
source");
+ }
+ }
+ }
+
private TLanceScanParams getOrCreateLanceScanParams() {
if (!params.isSetLanceScanParams()) {
params.setLanceScanParams(new TLanceScanParams());
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/datasource/tvf/source/TVFScanNode.java
b/fe/fe-core/src/main/java/org/apache/doris/datasource/tvf/source/TVFScanNode.java
index 99da0093513..2dccffa071b 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/datasource/tvf/source/TVFScanNode.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/datasource/tvf/source/TVFScanNode.java
@@ -18,7 +18,6 @@
package org.apache.doris.datasource.tvf.source;
import org.apache.doris.analysis.TupleDescriptor;
-import org.apache.doris.catalog.Env;
import org.apache.doris.catalog.FunctionGenTable;
import org.apache.doris.catalog.TableIf;
import org.apache.doris.common.DdlException;
@@ -33,13 +32,13 @@ import org.apache.doris.datasource.FileSplitter;
import org.apache.doris.datasource.TableFormatType;
import org.apache.doris.datasource.lance.LanceFragmentInfo;
import org.apache.doris.datasource.lance.LanceStorageOptions;
+import org.apache.doris.datasource.lance.source.LanceScanNode;
import org.apache.doris.datasource.lance.source.LanceSplit;
import org.apache.doris.planner.PlanNodeId;
import org.apache.doris.planner.ScanContext;
import org.apache.doris.qe.SessionVariable;
import org.apache.doris.spi.Split;
import org.apache.doris.statistics.StatisticalType;
-import org.apache.doris.system.Backend;
import org.apache.doris.tablefunction.ExternalFileTableValuedFunction;
import org.apache.doris.tablefunction.LocalTableValuedFunction;
import org.apache.doris.thrift.TBrokerFileStatus;
@@ -83,22 +82,25 @@ public class TVFScanNode extends FileQueryScanNode {
@Override
protected void initBackendPolicy() throws UserException {
- List<String> preferLocations = new ArrayList<>();
if (tableValuedFunction instanceof LocalTableValuedFunction) {
- // For local tvf, the backend was specified by backendId
- Long backendId = ((LocalTableValuedFunction)
tableValuedFunction).getBackendId();
+ long backendId =
+ ((LocalTableValuedFunction)
tableValuedFunction).getBackendIdForExecution();
if (backendId != -1) {
- // User has specified the backend, only use that backend
- // Otherwise, use all backends for shared storage.
- Backend backend =
Env.getCurrentSystemInfo().getBackend(backendId);
- if (backend == null) {
- throw new UserException("Backend " + backendId + " does
not exist");
- }
- preferLocations.add(backend.getHost());
+ backendPolicy.initWithBackendId(backendId);
+ numNodes = backendPolicy.numBackends();
+ return;
}
}
- backendPolicy.init(preferLocations);
+ backendPolicy.init();
numNodes = backendPolicy.numBackends();
+ if (tableValuedFunction.isLanceFormat()) {
+ boolean requiresCurrentReader = desc.getSlots().stream()
+ .anyMatch(slot -> slot.getColumn() != null
+ && tableValuedFunction.requiresCurrentLanceReader(
+ slot.getColumn().getName()));
+ LanceScanNode.checkAdditionalTypeBackendCompatibility(
+ requiresCurrentReader, backendPolicy.getBackends());
+ }
}
@Override
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/system/BeSelectionPolicy.java
b/fe/fe-core/src/main/java/org/apache/doris/system/BeSelectionPolicy.java
index 9699f8bab81..06a411b8694 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/system/BeSelectionPolicy.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/system/BeSelectionPolicy.java
@@ -62,6 +62,9 @@ public class BeSelectionPolicy {
public List<String> preferredLocations = new ArrayList<>();
+ // Empty means any backend ID is eligible.
+ public Set<Long> requiredBackendIds = Sets.newHashSet();
+
public boolean requireAliveBe = false;
private BeSelectionPolicy() {
@@ -131,6 +134,12 @@ public class BeSelectionPolicy {
return this;
}
+ /** Restrict candidates to the specified backend IDs. */
+ public Builder addRequiredBackendIds(Collection<Long> backendIds) {
+ policy.requiredBackendIds.addAll(backendIds);
+ return this;
+ }
+
public Builder setEnableRoundRobin(boolean enableRoundRobin) {
policy.enableRoundRobin = enableRoundRobin;
return this;
@@ -165,6 +174,7 @@ public class BeSelectionPolicy {
|| needLoadAvailable && !backend.isLoadAvailable()
|| needNonDecommissioned && (backend.isDecommissioned() ||
backend.isDecommissioning())
|| (!resourceTags.isEmpty() &&
!resourceTags.contains(backend.getLocationTag()))
+ || (!requiredBackendIds.isEmpty() &&
!requiredBackendIds.contains(backend.getId()))
|| storageMedium != null &&
!backend.hasSpecifiedStorageMedium(storageMedium)
|| (requireAliveBe && !backend.isAlive())) {
if (LOG.isDebugEnabled()) {
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/tablefunction/ExternalFileTableValuedFunction.java
b/fe/fe-core/src/main/java/org/apache/doris/tablefunction/ExternalFileTableValuedFunction.java
index 120156b54c7..3c4acfaab8a 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/tablefunction/ExternalFileTableValuedFunction.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/tablefunction/ExternalFileTableValuedFunction.java
@@ -142,6 +142,7 @@ public abstract class ExternalFileTableValuedFunction
extends TableValuedFunctio
private long tableId;
private long lanceDatasetVersion = -1;
private List<LanceFragmentInfo> lanceFragments = Collections.emptyList();
+ private Set<String> lanceCurrentReaderColumns = Collections.emptySet();
public abstract TFileType getTFileType();
@@ -169,6 +170,11 @@ public abstract class ExternalFileTableValuedFunction
extends TableValuedFunctio
return lanceFragments;
}
+ /** Returns whether a Lance column needs the current BE materialization
logic. */
+ public boolean requiresCurrentLanceReader(String columnName) {
+ return
lanceCurrentReaderColumns.contains(columnName.toLowerCase(Locale.ROOT));
+ }
+
public Map<String, String> getBackendConnectProperties() {
return backendConnectProperties;
}
@@ -385,11 +391,15 @@ public abstract class ExternalFileTableValuedFunction
extends TableValuedFunctio
List<Column> lanceColumns = new
ArrayList<>(metadata.getSchema().getFields().size());
Set<String> columnLowerNames = new HashSet<>();
+ Set<String> currentReaderColumns = new HashSet<>();
for (Field field : metadata.getSchema().getFields()) {
String lowerName = field.getName().toLowerCase(Locale.ROOT);
if (!columnLowerNames.add(lowerName)) {
throw new NotSupportedException("Repeated lowercase column
names: " + lowerName);
}
+ if (LanceTypeConverter.requiresCurrentBeReader(field)) {
+ currentReaderColumns.add(lowerName);
+ }
Type type;
try {
type = LanceTypeConverter.toDorisType(field);
@@ -403,6 +413,7 @@ public abstract class ExternalFileTableValuedFunction
extends TableValuedFunctio
lanceDatasetVersion = metadata.getVersion();
lanceFragments = Collections.unmodifiableList(fragments);
+ lanceCurrentReaderColumns =
Collections.unmodifiableSet(currentReaderColumns);
columns = lanceColumns;
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/tablefunction/LocalTableValuedFunction.java
b/fe/fe-core/src/main/java/org/apache/doris/tablefunction/LocalTableValuedFunction.java
index 17a6597c360..3e04d614aac 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/tablefunction/LocalTableValuedFunction.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/tablefunction/LocalTableValuedFunction.java
@@ -22,6 +22,7 @@ import org.apache.doris.analysis.StorageBackend.StorageType;
import org.apache.doris.catalog.Env;
import org.apache.doris.common.AnalysisException;
import org.apache.doris.common.UserException;
+import org.apache.doris.datasource.FederationBackendPolicy;
import org.apache.doris.datasource.property.storage.StorageProperties;
import org.apache.doris.proto.InternalService;
import org.apache.doris.proto.InternalService.PGlobResponse;
@@ -37,8 +38,6 @@ import com.google.common.base.Preconditions;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
-import java.util.Collections;
-import java.util.List;
import java.util.Map;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
@@ -95,24 +94,20 @@ public class LocalTableValuedFunction extends
ExternalFileTableValuedFunction {
}
}
+ /** Select the schema backend from the same eligibility policy used for
query execution. */
private void initializeBackendForRequest() throws AnalysisException {
- Backend be = null;
- if (backendId != -1) {
- be = Env.getCurrentSystemInfo().getBackend(backendId);
- backendIdForRequest = backendId;
- } else {
- Preconditions.checkState(sharedStorage);
- List<Long> beIds =
Env.getCurrentSystemInfo().getAllBackendByCurrentCluster(true);
- if (beIds.isEmpty()) {
- throw new AnalysisException("No available backend");
+ Preconditions.checkState(backendId != -1 || sharedStorage);
+ FederationBackendPolicy backendPolicy = new FederationBackendPolicy();
+ try {
+ if (backendId == -1) {
+ backendPolicy.init();
+ } else {
+ backendPolicy.initWithBackendId(backendId);
}
- Collections.shuffle(beIds);
- be = Env.getCurrentSystemInfo().getBackend(beIds.get(0));
- backendIdForRequest = be.getId();
- }
- if (be == null) {
- throw new AnalysisException("backend not found with backend_id = "
+ backendId);
+ } catch (UserException e) {
+ throw new AnalysisException("Failed to select local TVF backend: "
+ e.getMessage(), e);
}
+ backendIdForRequest = backendPolicy.getNextBe().getId();
}
private void getFileListFromBackend() throws AnalysisException {
@@ -164,6 +159,11 @@ public class LocalTableValuedFunction extends
ExternalFileTableValuedFunction {
return backendId;
}
+ /** Return the only backend that may execute this TVF, or -1 when
execution may be distributed. */
+ public long getBackendIdForExecution() {
+ return isLanceFormat() ? backendIdForRequest : backendId;
+ }
+
@Override
protected Backend getBackend() {
return Env.getCurrentSystemInfo().getBackend(backendIdForRequest);
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/datasource/lance/LanceTypeConverterTest.java
b/fe/fe-core/src/test/java/org/apache/doris/datasource/lance/LanceTypeConverterTest.java
index 54132a87756..9730b0af39e 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/datasource/lance/LanceTypeConverterTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/datasource/lance/LanceTypeConverterTest.java
@@ -102,16 +102,80 @@ public class LanceTypeConverterTest {
Field.nullable("uint64_col", new ArrowType.Int(64, false))));
}
+ /** Verifies Arrow Null and Duration mappings. */
@Test
- public void testExtensionAndDictionaryMarkersAreUnsupported() {
- Field extensionField = new Field(
- "json_col",
- new FieldType(
- true,
- ArrowType.Utf8.INSTANCE,
- null,
- Collections.singletonMap("ARROW:extension:name",
"lance.json")),
- Collections.emptyList());
+ public void testNullAndDurationMappings() {
+ Field nullField = Field.nullable("null_col", ArrowType.Null.INSTANCE);
+ Assertions.assertEquals(Type.NULL,
LanceTypeConverter.toDorisType(nullField));
+
Assertions.assertTrue(LanceTypeConverter.requiresCurrentBeReader(nullField));
+ for (TimeUnit unit : TimeUnit.values()) {
+ Field durationField =
+ Field.nullable("duration_col", new
ArrowType.Duration(unit));
+ Assertions.assertEquals(Type.BIGINT,
+ LanceTypeConverter.toDorisType(durationField));
+ Assertions.assertTrue(
+ LanceTypeConverter.requiresCurrentBeReader(durationField));
+ }
+ Field durationList = new Field(
+ "duration_list",
+ FieldType.nullable(ArrowType.List.INSTANCE),
+ Collections.singletonList(
+ Field.nullable("item", new
ArrowType.Duration(TimeUnit.MICROSECOND))));
+
Assertions.assertTrue(LanceTypeConverter.requiresCurrentBeReader(durationList));
+ }
+
+ /** Verifies nested Null fields remain unsupported. */
+ @Test
+ public void testNestedNullIsUnsupported() {
+ Field nullItem = Field.nullable("item", ArrowType.Null.INSTANCE);
+ Field nullList = new Field(
+ "null_list",
+ FieldType.nullable(ArrowType.List.INSTANCE),
+ Collections.singletonList(nullItem));
+ Field nullStruct = new Field(
+ "null_struct",
+ FieldType.nullable(ArrowType.Struct.INSTANCE),
+ Collections.singletonList(Field.nullable("value",
ArrowType.Null.INSTANCE)));
+
+ Assertions.assertEquals(Type.UNSUPPORTED,
LanceTypeConverter.toDorisType(nullList));
+ Assertions.assertEquals(Type.UNSUPPORTED,
LanceTypeConverter.toDorisType(nullStruct));
+ }
+
+ /** Verifies known extension mappings and storage validation. */
+ @Test
+ public void testKnownExtensionMappingsAndStorageValidation() {
+ Assertions.assertEquals(Type.JSONB, LanceTypeConverter.toDorisType(
+ extensionField("arrow_json_col", ArrowType.Utf8.INSTANCE,
"arrow.json")));
+ Assertions.assertTrue(LanceTypeConverter.requiresCurrentBeReader(
+ extensionField("arrow_json_col", ArrowType.Utf8.INSTANCE,
"arrow.json")));
+ Assertions.assertEquals(Type.JSONB, LanceTypeConverter.toDorisType(
+ extensionField("lance_json_col",
ArrowType.LargeBinary.INSTANCE, "lance.json")));
+
+ Field bfloat16Item = extensionField(
+ "item", new ArrowType.FixedSizeBinary(2), "lance.bfloat16");
+ Assertions.assertEquals(Type.FLOAT,
LanceTypeConverter.toDorisType(bfloat16Item));
+ Field bfloat16Vector = new Field(
+ "bfloat16_vector_col",
+ FieldType.nullable(new ArrowType.FixedSizeList(4)),
+ Collections.singletonList(bfloat16Item));
+ Assertions.assertEquals("array<float>",
+ LanceTypeConverter.toDorisType(bfloat16Vector).toSql());
+
Assertions.assertTrue(LanceTypeConverter.requiresCurrentBeReader(bfloat16Vector));
+ Assertions.assertFalse(LanceTypeConverter.requiresCurrentBeReader(
+ Field.nullable("ordinary", ArrowType.Utf8.INSTANCE)));
+
+ Assertions.assertEquals(Type.UNSUPPORTED,
LanceTypeConverter.toDorisType(
+ extensionField("invalid_json", ArrowType.Binary.INSTANCE,
"arrow.json")));
+ Assertions.assertEquals(Type.UNSUPPORTED,
LanceTypeConverter.toDorisType(
+ extensionField("invalid_bfloat16", new
ArrowType.FixedSizeBinary(4),
+ "lance.bfloat16")));
+ }
+
+ /** Verifies unknown extensions and dictionary fields remain unsupported.
*/
+ @Test
+ public void testUnknownExtensionAndDictionaryMarkersAreUnsupported() {
+ Field extensionField = extensionField(
+ "extension_col", ArrowType.Utf8.INSTANCE,
"doris.test.extension");
Assertions.assertEquals(Type.UNSUPPORTED,
LanceTypeConverter.toDorisType(extensionField));
Field dictionaryField = new Field(
@@ -132,19 +196,17 @@ public class LanceTypeConverterTest {
Collections.singletonMap("ARROW:extension:name",
"lance.blob.v2")),
Collections.emptyList());
Assertions.assertEquals(Type.UNSUPPORTED,
LanceTypeConverter.toDorisType(blobField));
+ }
- Field bfloat16Item = new Field(
- "item",
+ /** Creates a field with Arrow extension metadata. */
+ private static Field extensionField(String name, ArrowType storageType,
String extensionName) {
+ return new Field(
+ name,
new FieldType(
true,
- new ArrowType.FixedSizeBinary(2),
+ storageType,
null,
- Collections.singletonMap("ARROW:extension:name",
"lance.bfloat16")),
+ Collections.singletonMap("ARROW:extension:name",
extensionName)),
Collections.emptyList());
- Field bfloat16Vector = new Field(
- "bfloat16_vector_col",
- FieldType.nullable(new ArrowType.FixedSizeList(4)),
- Collections.singletonList(bfloat16Item));
- Assertions.assertEquals(Type.UNSUPPORTED,
LanceTypeConverter.toDorisType(bfloat16Vector));
}
}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/datasource/lance/source/LanceScanNodeTest.java
b/fe/fe-core/src/test/java/org/apache/doris/datasource/lance/source/LanceScanNodeTest.java
index 52f61bddca9..bdcf522c54a 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/datasource/lance/source/LanceScanNodeTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/datasource/lance/source/LanceScanNodeTest.java
@@ -27,6 +27,7 @@ import org.apache.doris.planner.PlanNodeId;
import org.apache.doris.planner.ScanContext;
import org.apache.doris.qe.SessionVariable;
import org.apache.doris.spi.Split;
+import org.apache.doris.system.Backend;
import org.apache.doris.thrift.TExternalSearchQuery;
import org.apache.doris.thrift.TExternalSearchRequest;
import org.apache.doris.thrift.TFileRangeDesc;
@@ -45,6 +46,7 @@ import org.apache.arrow.vector.types.pojo.Schema;
import org.junit.Assert;
import org.junit.Test;
import org.lance.index.IndexType;
+import org.mockito.Mockito;
import java.nio.ByteBuffer;
import java.util.Arrays;
@@ -54,6 +56,25 @@ import java.util.UUID;
public class LanceScanNodeTest {
+ // Verifies additional Lance encodings cannot run on a smooth-upgrade
source BE.
+ @Test
+ public void testAdditionalTypesRejectSmoothUpgradeSourceBackend() throws
Exception {
+ Backend currentBackend = Mockito.mock(Backend.class);
+ Backend smoothUpgradeSource = Mockito.mock(Backend.class);
+
Mockito.when(smoothUpgradeSource.isSmoothUpgradeSrc()).thenReturn(true);
+ Mockito.when(smoothUpgradeSource.getId()).thenReturn(10001L);
+
+ LanceScanNode.checkAdditionalTypeBackendCompatibility(
+ true, Arrays.asList(currentBackend));
+ LanceScanNode.checkAdditionalTypeBackendCompatibility(
+ false, Arrays.asList(currentBackend, smoothUpgradeSource));
+ UserException exception = Assert.assertThrows(UserException.class,
+ () -> LanceScanNode.checkAdditionalTypeBackendCompatibility(
+ true, Arrays.asList(currentBackend,
smoothUpgradeSource)));
+
+ Assert.assertTrue(exception.getMessage().contains("10001"));
+ }
+
@Test
public void testFragmentRowsDetermineSplitWeights() throws Exception {
LanceTableMetadata metadata = LanceTableMetadata.withoutIndexSegments(
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/datasource/tvf/source/TVFScanNodeTest.java
b/fe/fe-core/src/test/java/org/apache/doris/datasource/tvf/source/TVFScanNodeTest.java
index e079f3216f3..ef2dbb8f041 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/datasource/tvf/source/TVFScanNodeTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/datasource/tvf/source/TVFScanNodeTest.java
@@ -17,10 +17,16 @@
package org.apache.doris.datasource.tvf.source;
+import org.apache.doris.analysis.SlotDescriptor;
import org.apache.doris.analysis.SlotId;
import org.apache.doris.analysis.TupleDescriptor;
import org.apache.doris.analysis.TupleId;
+import org.apache.doris.catalog.Column;
import org.apache.doris.catalog.FunctionGenTable;
+import org.apache.doris.catalog.Type;
+import org.apache.doris.common.UserException;
+import org.apache.doris.datasource.ExternalScanNode;
+import org.apache.doris.datasource.FederationBackendPolicy;
import org.apache.doris.datasource.FileQueryScanNode;
import org.apache.doris.datasource.FileSplitter;
import org.apache.doris.datasource.lance.LanceFragmentInfo;
@@ -29,7 +35,9 @@ import org.apache.doris.planner.PlanNodeId;
import org.apache.doris.planner.ScanContext;
import org.apache.doris.qe.SessionVariable;
import org.apache.doris.spi.Split;
+import org.apache.doris.system.Backend;
import org.apache.doris.tablefunction.ExternalFileTableValuedFunction;
+import org.apache.doris.tablefunction.LocalTableValuedFunction;
import org.apache.doris.thrift.TBrokerFileStatus;
import org.apache.doris.thrift.TFileFormatType;
import org.apache.doris.thrift.TFileRangeDesc;
@@ -41,6 +49,7 @@ import org.junit.Assert;
import org.junit.Test;
import org.mockito.Mockito;
+import java.lang.reflect.Field;
import java.lang.reflect.Method;
import java.util.Arrays;
import java.util.Collections;
@@ -204,4 +213,61 @@ public class TVFScanNodeTest {
Assert.assertEquals(0L,
range.getTableFormatParams().getLanceParams().getVersion());
Assert.assertFalse(range.getTableFormatParams().getLanceParams().isSetFragmentIds());
}
+
+ // Verifies local Lance execution is restricted to the backend that
provided its schema.
+ @Test
+ public void testLocalLancePinsExecutionToSchemaBackend() throws Exception {
+ SessionVariable sv = new SessionVariable();
+ TupleDescriptor desc = new TupleDescriptor(new TupleId(0));
+ FunctionGenTable table = Mockito.mock(FunctionGenTable.class);
+ LocalTableValuedFunction tvf =
Mockito.mock(LocalTableValuedFunction.class);
+ Mockito.when(table.getTvf()).thenReturn(tvf);
+ Mockito.when(tvf.getBackendIdForExecution()).thenReturn(101L);
+ desc.setTable(table);
+
+ TVFScanNode node = new TVFScanNode(new PlanNodeId(0), desc, false, sv,
ScanContext.EMPTY);
+ FederationBackendPolicy backendPolicy =
Mockito.mock(FederationBackendPolicy.class);
+ Mockito.when(backendPolicy.numBackends()).thenReturn(1);
+ Field backendPolicyField =
ExternalScanNode.class.getDeclaredField("backendPolicy");
+ backendPolicyField.setAccessible(true);
+ backendPolicyField.set(node, backendPolicy);
+
+ node.initBackendPolicy();
+
+ Mockito.verify(backendPolicy).initWithBackendId(101L);
+ Mockito.verify(backendPolicy, Mockito.never()).init();
+ }
+
+ // Verifies S3 Lance projections reject a smooth-upgrade source BE.
+ @Test
+ public void testS3LanceAdditionalTypesRejectSmoothUpgradeSource() throws
Exception {
+ SessionVariable sv = new SessionVariable();
+ TupleDescriptor desc = new TupleDescriptor(new TupleId(0));
+ SlotDescriptor slot = new SlotDescriptor(new SlotId(1), desc);
+ slot.setColumn(new Column("json_value", Type.JSONB));
+ desc.addSlot(slot);
+ FunctionGenTable table = Mockito.mock(FunctionGenTable.class);
+ ExternalFileTableValuedFunction tvf =
+ Mockito.mock(ExternalFileTableValuedFunction.class);
+ Mockito.when(table.getTvf()).thenReturn(tvf);
+ Mockito.when(tvf.isLanceFormat()).thenReturn(true);
+
Mockito.when(tvf.requiresCurrentLanceReader("json_value")).thenReturn(true);
+ desc.setTable(table);
+
+ Backend smoothUpgradeSource = Mockito.mock(Backend.class);
+
Mockito.when(smoothUpgradeSource.isSmoothUpgradeSrc()).thenReturn(true);
+ Mockito.when(smoothUpgradeSource.getId()).thenReturn(102L);
+ FederationBackendPolicy backendPolicy =
Mockito.mock(FederationBackendPolicy.class);
+ Mockito.when(backendPolicy.getBackends())
+ .thenReturn(Collections.singletonList(smoothUpgradeSource));
+ TVFScanNode node = new TVFScanNode(new PlanNodeId(0), desc, false, sv,
ScanContext.EMPTY);
+ Field backendPolicyField =
ExternalScanNode.class.getDeclaredField("backendPolicy");
+ backendPolicyField.setAccessible(true);
+ backendPolicyField.set(node, backendPolicy);
+
+ UserException exception =
+ Assert.assertThrows(UserException.class,
node::initBackendPolicy);
+
+ Assert.assertTrue(exception.getMessage().contains("102"));
+ }
}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/system/SystemInfoServiceTest.java
b/fe/fe-core/src/test/java/org/apache/doris/system/SystemInfoServiceTest.java
index f79a5a7020f..a7d479e4184 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/system/SystemInfoServiceTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/system/SystemInfoServiceTest.java
@@ -36,6 +36,7 @@ import com.google.common.collect.Sets;
import org.junit.Assert;
import org.junit.Before;
import org.junit.Test;
+import org.mockito.Mockito;
import java.io.DataInputStream;
import java.io.DataOutputStream;
@@ -379,6 +380,29 @@ public class SystemInfoServiceTest {
Assert.assertEquals(1, infoService.selectBackendIdsByPolicy(policy4,
1).size());
}
+ // Verifies an exact backend-ID constraint cannot fall back to another
eligible backend.
+ @Test
+ public void testRequiredBackendIdsSelect() {
+ Backend selected = Mockito.mock(Backend.class);
+ Backend excluded = Mockito.mock(Backend.class);
+ Mockito.when(selected.getId()).thenReturn(10001L);
+ Mockito.when(excluded.getId()).thenReturn(10002L);
+ Mockito.when(selected.isQueryAvailable()).thenReturn(true);
+ Mockito.when(excluded.isQueryAvailable()).thenReturn(true);
+ BeSelectionPolicy policy = new BeSelectionPolicy.Builder()
+
.addRequiredBackendIds(Collections.singletonList(selected.getId()))
+ .build();
+
+ List<Backend> candidates =
+ policy.getCandidateBackends(Lists.newArrayList(selected,
excluded));
+
+ Assert.assertEquals(Collections.singletonList(selected), candidates);
+
+ Mockito.when(selected.isComputeNode()).thenReturn(true);
+ Assert.assertTrue(
+ policy.getCandidateBackends(Lists.newArrayList(selected,
excluded)).isEmpty());
+ }
+
@Test
public void testSelectBackendIdsForReplicaCreation() throws Exception {
addBackend(10001, "192.168.1.1", 9050);
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/tablefunction/ExternalFileTableValuedFunctionTest.java
b/fe/fe-core/src/test/java/org/apache/doris/tablefunction/ExternalFileTableValuedFunctionTest.java
index e9741855522..c30fc885b94 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/tablefunction/ExternalFileTableValuedFunctionTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/tablefunction/ExternalFileTableValuedFunctionTest.java
@@ -23,15 +23,24 @@ import org.apache.doris.common.AnalysisException;
import org.apache.doris.common.Config;
import org.apache.doris.common.util.FileFormatConstants;
import org.apache.doris.common.util.FileFormatUtils;
+import org.apache.doris.datasource.lance.LanceTableMetadata;
import org.apache.doris.datasource.property.fileformat.FileFormatProperties;
import
org.apache.doris.datasource.property.fileformat.LanceFileFormatProperties;
import org.apache.doris.thrift.TFileFormatType;
import com.google.common.collect.Lists;
import com.google.common.collect.Maps;
+import org.apache.arrow.vector.types.TimeUnit;
+import org.apache.arrow.vector.types.pojo.ArrowType;
+import org.apache.arrow.vector.types.pojo.Field;
+import org.apache.arrow.vector.types.pojo.FieldType;
+import org.apache.arrow.vector.types.pojo.Schema;
import org.junit.Assert;
import org.junit.Test;
+import org.mockito.Mockito;
+import java.util.Arrays;
+import java.util.Collections;
import java.util.List;
import java.util.Map;
@@ -126,4 +135,55 @@ public class ExternalFileTableValuedFunctionTest {
Assert.fail();
}
}
+
+ // Verifies a shared-storage Lance TVF executes on the backend that
provided its schema.
+ @Test
+ public void testLocalLanceExecutionUsesSchemaBackend() throws Exception {
+ LocalTableValuedFunction tvf =
+ Mockito.mock(LocalTableValuedFunction.class,
Mockito.CALLS_REAL_METHODS);
+ setLongField(tvf, "backendId", -1L);
+ setLongField(tvf, "backendIdForRequest", 23L);
+
+ Mockito.doReturn(true).when(tvf).isLanceFormat();
+ Assert.assertEquals(23L, tvf.getBackendIdForExecution());
+
+ Mockito.doReturn(false).when(tvf).isLanceFormat();
+ Assert.assertEquals(-1L, tvf.getBackendIdForExecution());
+ }
+
+ // Verifies S3 Lance metadata records which columns require the current BE
reader.
+ @Test
+ public void testLanceMetadataTracksCurrentReaderColumns() throws Exception
{
+ ExternalFileTableValuedFunction tvf =
+ Mockito.mock(ExternalFileTableValuedFunction.class,
Mockito.CALLS_REAL_METHODS);
+ Field jsonField = new Field(
+ "json_value",
+ new FieldType(true, ArrowType.Utf8.INSTANCE, null,
+ Collections.singletonMap("ARROW:extension:name",
"arrow.json")),
+ Collections.emptyList());
+ LanceTableMetadata metadata = LanceTableMetadata.withoutIndexSegments(
+ "s3://bucket/table.lance", 1L,
+ new Schema(Arrays.asList(
+ jsonField,
+ Field.nullable("null_value", ArrowType.Null.INSTANCE),
+ Field.nullable("duration_value",
+ new ArrowType.Duration(TimeUnit.MILLISECOND)),
+ Field.nullable("ordinary", ArrowType.Utf8.INSTANCE))),
+ Collections.emptyList(), Collections.emptyMap());
+
+ tvf.setLanceTableMetadata(metadata);
+
+ Assert.assertTrue(tvf.requiresCurrentLanceReader("JSON_VALUE"));
+ Assert.assertTrue(tvf.requiresCurrentLanceReader("null_value"));
+ Assert.assertTrue(tvf.requiresCurrentLanceReader("DURATION_VALUE"));
+ Assert.assertFalse(tvf.requiresCurrentLanceReader("ordinary"));
+ }
+
+ // Sets a private long field without invoking the table function's
environment-dependent constructor.
+ private static void setLongField(Object target, String fieldName, long
value) throws Exception {
+ java.lang.reflect.Field field =
+ LocalTableValuedFunction.class.getDeclaredField(fieldName);
+ field.setAccessible(true);
+ field.setLong(target, value);
+ }
}
diff --git
a/regression-test/data/external_table_p0/lance/test_lance_catalog_all_types.out
b/regression-test/data/external_table_p0/lance/test_lance_catalog_all_types.out
index 3f4edf5f574..78920970c9b 100644
---
a/regression-test/data/external_table_p0/lance/test_lance_catalog_all_types.out
+++
b/regression-test/data/external_table_p0/lance/test_lance_catalog_all_types.out
@@ -1,6 +1,6 @@
-- This file is automatically generated. You should know what you did if you
want to edit this
-- !sql_1 --
-bfloat16_vector_col unknown type: UNSUPPORTED_TYPE Yes false \N
+bfloat16_vector_col array<float> Yes false \N
binary_col varbinary(2147483647) Yes false \N
blob_col unknown type: UNSUPPORTED_TYPE Yes false \N
bool_col boolean Yes false \N
@@ -9,10 +9,10 @@ date64_col date Yes false \N
decimal128_col decimal(38,10) Yes false \N
decimal256_col decimal(76,38) Yes false \N
dictionary_col smallint Yes false \N
-duration_ms_col unknown type: UNSUPPORTED_TYPE Yes false \N
-duration_ns_col unknown type: UNSUPPORTED_TYPE Yes false \N
-duration_s_col unknown type: UNSUPPORTED_TYPE Yes false \N
-duration_us_col unknown type: UNSUPPORTED_TYPE Yes false \N
+duration_ms_col bigint Yes false \N
+duration_ns_col bigint Yes false \N
+duration_s_col bigint Yes false \N
+duration_us_col bigint Yes false \N
fixed_size_binary_col varbinary(16) Yes false \N
fixed_size_list_float16_col array<float> Yes false \N
fixed_size_list_float32_col array<float> Yes false \N
@@ -26,7 +26,7 @@ int16_col smallint Yes false \N
int32_col int Yes false \N
int64_col bigint Yes false \N
int8_col tinyint Yes false \N
-json_col unknown type: UNSUPPORTED_TYPE Yes false \N
+json_col json Yes false \N
large_binary_col varbinary(2147483647) Yes false \N
large_list_col array<float> Yes false \N
large_list_struct_col array<struct<name:text,value:int>> Yes false
\N
@@ -34,7 +34,7 @@ large_utf8_col text Yes false \N
list_col array<int> Yes false \N
list_struct_col array<struct<name:text,value:int>> Yes false
\N
map_col map<text,int> Yes false \N
-null_col unknown type: UNSUPPORTED_TYPE Yes false \N
+null_col null_type Yes false \N
row_id bigint Yes false \N
struct_col struct<name:text,score:double> Yes false \N
time32_ms_col time(3) Yes false \N
@@ -55,3 +55,5 @@ utf8_col text Yes false \N
-- !select --
true -8 8 -16 16 -32 32 -64 64 1.5
2.5 3.5 1234567890.1234567890
12345678901234567890123456789012345678.12345678901234567890123456789012345678
Doris and Lance large string 0001FF 6C617267652062696E617279
30313233343536373839616263646566 2026-07-28 2026-07-28
12:34:56 12:34:56.123 12:34:56.123456 12:34:56.123456
2026-07-28T12:34:56 2026-07-28T12:34:56.123 2026-07-28T12:34:56.123456
2026-07-28T12:34:56.123456 2026-07-28 20:34:56.123456+08:00
{"name":"doris", "score":1} [1, 2, 3] [{"name":" [...]
+-- !additional_lance_types --
+true 1 1000 1000000 1000000000
{"engine":"doris","format":"lance"} [1, 2, 3, 4]
diff --git a/regression-test/data/external_table_p0/lance/test_lance_s3_tvf.out
b/regression-test/data/external_table_p0/lance/test_lance_s3_tvf.out
index e689a8a9de3..5b81da80604 100644
--- a/regression-test/data/external_table_p0/lance/test_lance_s3_tvf.out
+++ b/regression-test/data/external_table_p0/lance/test_lance_s3_tvf.out
@@ -1,6 +1,6 @@
-- This file is automatically generated. You should know what you did if you
want to edit this
-- !desc --
-bfloat16_vector_col unknown type: UNSUPPORTED_TYPE Yes false \N
NONE
+bfloat16_vector_col array<float> Yes false \N NONE
binary_col varbinary(2147483647) Yes false \N NONE
blob_col unknown type: UNSUPPORTED_TYPE Yes false \N NONE
bool_col boolean Yes false \N NONE
@@ -9,10 +9,10 @@ date64_col date Yes false \N NONE
decimal128_col decimal(38,10) Yes false \N NONE
decimal256_col decimal(76,38) Yes false \N NONE
dictionary_col smallint Yes false \N NONE
-duration_ms_col unknown type: UNSUPPORTED_TYPE Yes false \N
NONE
-duration_ns_col unknown type: UNSUPPORTED_TYPE Yes false \N
NONE
-duration_s_col unknown type: UNSUPPORTED_TYPE Yes false \N NONE
-duration_us_col unknown type: UNSUPPORTED_TYPE Yes false \N
NONE
+duration_ms_col bigint Yes false \N NONE
+duration_ns_col bigint Yes false \N NONE
+duration_s_col bigint Yes false \N NONE
+duration_us_col bigint Yes false \N NONE
fixed_size_binary_col varbinary(16) Yes false \N NONE
fixed_size_list_float16_col array<float> Yes false \N NONE
fixed_size_list_float32_col array<float> Yes false \N NONE
@@ -26,7 +26,7 @@ int16_col smallint Yes false \N NONE
int32_col int Yes false \N NONE
int64_col bigint Yes false \N NONE
int8_col tinyint Yes false \N NONE
-json_col unknown type: UNSUPPORTED_TYPE Yes false \N NONE
+json_col json Yes false \N NONE
large_binary_col varbinary(2147483647) Yes false \N NONE
large_list_col array<float> Yes false \N NONE
large_list_struct_col array<struct<name:text,value:int>> Yes false
\N NONE
@@ -34,7 +34,7 @@ large_utf8_col text Yes false \N NONE
list_col array<int> Yes false \N NONE
list_struct_col array<struct<name:text,value:int>> Yes false
\N NONE
map_col map<text,int> Yes false \N NONE
-null_col unknown type: UNSUPPORTED_TYPE Yes false \N NONE
+null_col null_type Yes false \N NONE
row_id bigint Yes false \N NONE
struct_col struct<name:text,score:double> Yes false \N NONE
time32_ms_col time(3) Yes false \N NONE
@@ -64,3 +64,6 @@ utf8_col text Yes false \N NONE
-- !representative_types --
true 1.5
12345678901234567890123456789012345678.12345678901234567890123456789012345678
0001FF 12:34:56.123 12:34:56.123456 2026-07-28T12:34:56.123456
2026-07-28 20:34:56.123456+08:00 [0, 1, 2, 3] {"doris":1, "lance":2}
+
+-- !additional_types --
+true 1 1000 1000000 1000000000
{"engine":"doris","format":"lance"} [1, 2, 3, 4]
diff --git
a/regression-test/suites/external_table_p0/lance/test_lance_catalog_all_types.groovy
b/regression-test/suites/external_table_p0/lance/test_lance_catalog_all_types.groovy
index f87e2d70b5c..fec04ea6325 100644
---
a/regression-test/suites/external_table_p0/lance/test_lance_catalog_all_types.groovy
+++
b/regression-test/suites/external_table_p0/lance/test_lance_catalog_all_types.groovy
@@ -29,7 +29,9 @@ suite("test_lance_catalog_all_types","p0,external") {
String catalogName = "test_lance_catalog_all_types"
String databaseName = "default"
String tableName = "all_types"
- String unsupported = "unknown type: UNSUPPORTED_TYPE"
+ def scannerV2Rows = sql """SHOW VARIABLES LIKE 'enable_file_scanner_v2'"""
+ assertEquals(1, scannerV2Rows.size())
+ String originalScannerV2 = scannerV2Rows[0][1].toString()
/*
* Lance/Arrow to Doris type mapping exercised by this fixture:
@@ -38,7 +40,7 @@ suite("test_lance_catalog_all_types","p0,external") {
*
* | Lance / Arrow type | Doris DESC type | Status / notes |
* |---|---|---|
- * | null | UNSUPPORTED | Unsupported |
+ * | null | null_type | Supported; every value is SQL NULL |
* | bool | boolean | Supported |
* | int8 | tinyint | Supported |
* | uint8 | smallint | Supported by lossless widening |
@@ -65,7 +67,7 @@ suite("test_lance_catalog_all_types","p0,external") {
* | timestamp(ms) | datetime(3) | Supported |
* | timestamp(us/ns) | datetime(6) | Nanoseconds are truncated to
microseconds |
* | timestamp(us, UTC) | timestamptz(6) | UTC instant; rendered in the
Doris session timezone |
- * | duration(s/ms/us/ns) | UNSUPPORTED | Unsupported |
+ * | duration(s/ms/us/ns) | bigint | Exact signed count in the Arrow
field's declared unit |
* | struct | struct<...> | Child types are converted recursively |
* | list / large_list / fixed_size_list | array<...> | Item type is
converted recursively |
* | fixed_size_list<uint8> | array<smallint> | Item type is widened
recursively |
@@ -73,8 +75,8 @@ suite("test_lance_catalog_all_types","p0,external") {
* | map | map<...,...> | Key/value types are converted recursively |
* | dictionary<int16, utf8> | smallint | SDK currently exposes only the
physical index type |
* | Lance Blob v2 | UNSUPPORTED | Blob materialization is not implemented
|
- * | Arrow JSON extension | UNSUPPORTED | JSON extension decoding is not
implemented |
- * | fixed_size_list<Lance BFloat16> | UNSUPPORTED | BFloat16 decoding is
not implemented |
+ * | Arrow JSON extension | json | Returned as Doris JSON |
+ * | fixed_size_list<Lance BFloat16> | array<float> | BFloat16 values are
widened exactly |
*
* Dataset.getSchema() preserves extension metadata for Blob, JSON, and
BFloat16,
* allowing Doris to reject them without interpreting their physical
storage.
@@ -85,6 +87,7 @@ suite("test_lance_catalog_all_types","p0,external") {
sql """DROP CATALOG IF EXISTS `${catalogName}`"""
try {
+ sql """SET enable_file_scanner_v2 = true"""
sql """
CREATE CATALOG `${catalogName}` PROPERTIES (
"type" = "lance",
@@ -126,6 +129,20 @@ suite("test_lance_catalog_all_types","p0,external") {
where row_id = 1;
"""
+ // Verify both schema mappings and values for the additional Lance
types.
+ qt_additional_lance_types """
+ SELECT
+ null_col IS NULL AS null_is_null,
+ duration_s_col,
+ duration_ms_col,
+ duration_us_col,
+ duration_ns_col,
+ CAST(json_col AS STRING) AS json_col,
+ bfloat16_vector_col
+ FROM `${catalogName}`.`${databaseName}`.`${tableName}`
+ WHERE row_id = 1
+ """
+
def timestampRows = sql """
SELECT timestamp_us_col, timestamp_us_utc_col
FROM `${catalogName}`.`${databaseName}`.`${tableName}`
@@ -145,6 +162,7 @@ suite("test_lance_catalog_all_types","p0,external") {
assertEquals("2026-07-28 12:34:56.123456+00:00",
timestampRows[0][1].toString())
} finally {
+ sql """SET enable_file_scanner_v2 = ${originalScannerV2}"""
// sql """DROP CATALOG IF EXISTS `${catalogName}`"""
}
}
diff --git
a/regression-test/suites/external_table_p0/lance/test_lance_s3_tvf.groovy
b/regression-test/suites/external_table_p0/lance/test_lance_s3_tvf.groovy
index f6a3aa5e734..3509690b3e5 100644
--- a/regression-test/suites/external_table_p0/lance/test_lance_s3_tvf.groovy
+++ b/regression-test/suites/external_table_p0/lance/test_lance_s3_tvf.groovy
@@ -80,10 +80,18 @@ suite("test_lance_s3_tvf", "p0,external") {
WHERE row_id = 1
"""
- test {
- sql """SELECT null_col FROM ${lanceTvf}"""
- exception "is unsupported for Nereids"
- }
+ qt_additional_types """
+ SELECT
+ null_col IS NULL AS null_is_null,
+ duration_s_col,
+ duration_ms_col,
+ duration_us_col,
+ duration_ns_col,
+ CAST(json_col AS STRING) AS json_col,
+ bfloat16_vector_col
+ FROM ${lanceTvf}
+ WHERE row_id = 1
+ """
} finally {
sql """SET enable_file_scanner_v2 = ${originalScannerV2}"""
sql """SET time_zone = '${originalTimeZone}'"""
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]