This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new c1bab9e8cf8 [refactor](be) move orc related serde method to specific
datatypes (#66515)
c1bab9e8cf8 is described below
commit c1bab9e8cf80fc5e0c52914e9d4aa69717ea2e7a
Author: yiguolei <[email protected]>
AuthorDate: Fri Aug 7 09:13:57 2026 +0800
[refactor](be) move orc related serde method to specific datatypes (#66515)
### What problem does this PR solve?
Issue Number: close #xxx
Related PR: #xxx
Problem Summary:
### Release note
None
### Check List (For Author)
- Test <!-- At least one of them must be included. -->
- [ ] Regression test
- [ ] Unit Test
- [ ] Manual test (add detailed scripts or steps below)
- [ ] No need to test or manual test. Explain why:
- [ ] This is a refactor/code format and no logic has been changed.
- [ ] Previous test can cover this change.
- [ ] No code files have been changed.
- [ ] Other reason <!-- Add your reason? -->
- Behavior changed:
- [ ] No.
- [ ] Yes. <!-- Explain the behavior change -->
- Does this need documentation?
- [ ] No.
- [ ] Yes. <!-- Add document PR link here. eg:
https://github.com/apache/doris-website/pull/1214 -->
### Check List (For Reviewer who merge this PR)
- [ ] Confirm the release note
- [ ] Confirm test cases
- [ ] Confirm document
- [ ] Add branch pick label <!-- Add branch pick label that this PR
should merge into -->
---
.../core/data_type_serde/data_type_array_serde.cpp | 54 +-
.../data_type_serde/data_type_datetimev2_serde.cpp | 61 ++
.../data_type_serde/data_type_datev2_serde.cpp | 43 +
.../data_type_serde/data_type_decimal_serde.cpp | 122 +++
.../core/data_type_serde/data_type_map_serde.cpp | 68 +-
.../data_type_serde/data_type_nullable_serde.cpp | 47 ++
.../data_type_serde/data_type_number_serde.cpp | 128 +++
be/src/core/data_type_serde/data_type_serde.cpp | 901 ---------------------
.../data_type_serde/data_type_string_serde.cpp | 99 +++
.../data_type_serde/data_type_struct_serde.cpp | 64 ++
.../data_type_timestamptz_serde.cpp | 51 ++
be/src/core/data_type_serde/orc_serde_utils.cpp | 161 ++++
be/src/core/data_type_serde/orc_serde_utils.h | 40 +
.../format/transformer/vorc_transformer_test.cpp | 27 +-
14 files changed, 959 insertions(+), 907 deletions(-)
diff --git a/be/src/core/data_type_serde/data_type_array_serde.cpp
b/be/src/core/data_type_serde/data_type_array_serde.cpp
index b63f081e594..bfab0f1eddb 100644
--- a/be/src/core/data_type_serde/data_type_array_serde.cpp
+++ b/be/src/core/data_type_serde/data_type_array_serde.cpp
@@ -404,7 +404,7 @@ Status DataTypeArraySerDe::write_column_to_orc(const
std::string& timezone, cons
packed_nested_size,
arena, options));
// String batches borrow their source bytes, but the packed column is
local to this call;
// keep only those borrowed leaves in the write Arena until Writer::add()
consumes them.
- copy_orc_string_data_to_arena(cur_batch->elements.get(), arena);
+ orc_serde_utils::copy_orc_string_data_to_arena(cur_batch->elements.get(),
arena);
cur_batch->elements->numElements = packed_nested_size;
cur_batch->numElements = end - start;
@@ -627,4 +627,56 @@ bool DataTypeArraySerDe::write_column_to_hive_text(const
IColumn& column, Buffer
return true;
}
+namespace {
+
+Status decode_list_orc_values(const DataTypeSerDeSPtr& nested_serde, IColumn&
nested_column,
+ const OrcDecodedColumnView& orc_view) {
+ const auto* orc_list = dynamic_cast<const
::orc::ListVectorBatch*>(orc_view.batch);
+ if (orc_list == nullptr) {
+ return Status::InternalError("Unexpected ORC list batch type {}",
+ orc_view.batch->toString());
+ }
+ DORIS_CHECK(orc_view.file_type != nullptr);
+ DORIS_CHECK(orc_view.selected_type != nullptr);
+ DORIS_CHECK(orc_view.file_type->getSubtypeCount() == 1);
+ DORIS_CHECK(orc_view.selected_type->getSubtypeCount() == 1);
+ DORIS_CHECK(orc_list->elements != nullptr);
+ const auto* file_element_type = orc_view.file_type->getSubtype(0);
+ const auto* selected_element_type = orc_view.selected_type->getSubtype(0);
+ DORIS_CHECK(file_element_type != nullptr);
+ DORIS_CHECK(selected_element_type != nullptr);
+
+ auto& array_column = assert_cast<ColumnArray&>(nested_column);
+ size_t element_size = 0;
+ std::vector<size_t> element_selection;
+ RETURN_IF_ERROR(orc_serde_utils::append_orc_offsets(
+ array_column.get_offsets(), orc_list->offsets, orc_view.rows,
&element_size,
+ orc_view.selected_rows, &element_selection));
+ auto element_column = array_column.get_data_ptr()->assert_mutable();
+ const auto child_rows = orc_view.selected_rows == nullptr
+ ? element_size
+ :
static_cast<size_t>(orc_list->elements->numElements);
+ const auto* child_selection = orc_view.selected_rows == nullptr ? nullptr
: &element_selection;
+ auto child_view = orc_serde_utils::make_child_orc_view(
+ orc_view, file_element_type, selected_element_type,
orc_list->elements.get(),
+ child_rows, child_selection);
+ RETURN_IF_ERROR(
+ orc_serde_utils::read_orc_child_column(nested_serde,
element_column, child_view));
+ array_column.get_data_ptr() = std::move(element_column);
+ return Status::OK();
+}
+
+} // namespace
+
+Status DataTypeArraySerDe::read_column_from_orc(IColumn& column,
+ const OrcDecodedColumnView&
view) const {
+ DORIS_CHECK(view.file_type != nullptr);
+ DORIS_CHECK(view.batch != nullptr);
+ DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::LIST);
+ if (orc_serde_utils::orc_decode_row_count(view.rows, view.selected_rows)
== 0) {
+ return Status::OK();
+ }
+ return decode_list_orc_values(nested_serde, column, view);
+}
+
} // namespace doris
diff --git a/be/src/core/data_type_serde/data_type_datetimev2_serde.cpp
b/be/src/core/data_type_serde/data_type_datetimev2_serde.cpp
index 7b5ad093916..78370a858c9 100644
--- a/be/src/core/data_type_serde/data_type_datetimev2_serde.cpp
+++ b/be/src/core/data_type_serde/data_type_datetimev2_serde.cpp
@@ -31,6 +31,7 @@
#include "core/data_type/primitive_type.h"
#include "core/data_type_serde/arrow_validation.h"
#include "core/data_type_serde/decoded_column_view.h"
+#include "core/data_type_serde/orc_serde_utils.h"
#include "core/data_type_serde/parquet_decode_source.h"
#include "core/data_type_serde/parquet_timestamp.h"
#include "core/types.h"
@@ -51,6 +52,52 @@ static const int64_t micro_to_nano_second = 1000;
namespace {
+Status decode_timestamp_orc_values(IColumn& nested_column, const
OrcDecodedColumnView& orc_view,
+ const cctz::time_zone& timezone) {
+ const auto* orc_batch = dynamic_cast<const
::orc::TimestampVectorBatch*>(orc_view.batch);
+ if (orc_batch == nullptr) {
+ return Status::InternalError("Unexpected ORC timestamp batch type {}",
+ orc_view.batch->toString());
+ }
+ auto& data = assert_cast<ColumnDateTimeV2&>(nested_column).get_data();
+ const size_t old_data_size = data.size();
+ const auto output_rows =
+ orc_serde_utils::orc_decode_row_count(orc_view.rows,
orc_view.selected_rows);
+ data.resize(old_data_size + output_rows);
+ for (size_t row = 0; row < output_rows; ++row) {
+ const auto source_row = orc_serde_utils::orc_source_row_at(row,
orc_view.selected_rows);
+ if (orc_serde_utils::orc_row_is_null(*orc_view.batch, source_row)) {
+ data[old_data_size + row] = DateV2Value<DateTimeV2ValueType> {};
+ continue;
+ }
+ auto& value =
+
reinterpret_cast<DateV2Value<DateTimeV2ValueType>&>(data[old_data_size + row]);
+ orc_serde_utils::RoundedOrcTimestamp timestamp;
+ auto status = orc_serde_utils::round_orc_timestamp_to_microseconds(
+ orc_batch->data[source_row],
orc_batch->nanoseconds[source_row], ×tamp);
+ if (!status.ok()) {
+ data.resize(old_data_size);
+ return status;
+ }
+ value.from_unixtime(orc_batch->data[source_row], timezone);
+ if (!value.is_valid_date()) {
+ data.resize(old_data_size);
+ return Status::DataQualityError(
+ "Decoded ORC timestamp is outside the target timezone
range");
+ }
+ value.set_microsecond(timestamp.microseconds);
+ // Plain ORC TIMESTAMP is a civil value. Carry after timezone
conversion so a fractional
+ // round does not jump backward or skip an hour at a daylight-saving
transition.
+ if (timestamp.carry &&
+ !value.date_add_interval<TimeUnit::SECOND>(TimeInterval
{TimeUnit::SECOND, 1, false})) {
+ data.resize(old_data_size);
+ return Status::DataQualityError(
+ "Decoded ORC timestamp is outside the target timezone
range");
+ }
+ }
+ return Status::OK();
+}
+
Status append_datetimev2_from_epoch_micros(ColumnDateTimeV2::Container& data,
int64_t timestamp_micros) {
static constexpr int64_t MICROS_PER_SECOND = 1000000;
@@ -926,4 +973,18 @@ template Status
DataTypeDateTimeV2SerDe::from_decimal_strict_mode_batch<DataType
const DataTypeDecimal128::ColumnType& decimal_col, IColumn&
target_col) const;
template Status
DataTypeDateTimeV2SerDe::from_decimal_strict_mode_batch<DataTypeDecimal256>(
const DataTypeDecimal256::ColumnType& decimal_col, IColumn&
target_col) const;
+
+Status DataTypeDateTimeV2SerDe::read_column_from_orc(IColumn& column,
+ const
OrcDecodedColumnView& view) const {
+ DORIS_CHECK(view.file_type != nullptr);
+ DORIS_CHECK(view.batch != nullptr);
+ const auto kind = view.file_type->getKind();
+ DORIS_CHECK(kind == ::orc::TypeKind::TIMESTAMP || kind ==
::orc::TypeKind::TIMESTAMP_INSTANT);
+ DORIS_CHECK(view.timezone != nullptr);
+ if (orc_serde_utils::orc_decode_row_count(view.rows, view.selected_rows)
== 0) {
+ return Status::OK();
+ }
+ return decode_timestamp_orc_values(column, view, *view.timezone);
+}
+
} // namespace doris
diff --git a/be/src/core/data_type_serde/data_type_datev2_serde.cpp
b/be/src/core/data_type_serde/data_type_datev2_serde.cpp
index 7b6b57fbd08..be9c841d233 100644
--- a/be/src/core/data_type_serde/data_type_datev2_serde.cpp
+++ b/be/src/core/data_type_serde/data_type_datev2_serde.cpp
@@ -22,6 +22,7 @@
#include <fmt/core.h>
#include <cstdint>
+#include <vector>
#include "common/config.h"
#include "core/column/column_const.h"
@@ -30,11 +31,13 @@
#include "core/data_type/define_primitive_type.h"
#include "core/data_type_serde/arrow_validation.h"
#include "core/data_type_serde/decoded_column_view.h"
+#include "core/data_type_serde/orc_serde_utils.h"
#include "core/data_type_serde/parquet_decode_source.h"
#include "core/types.h"
#include "core/value/vdatetime_value.h"
#include "exprs/function/cast/cast_to_datev2_impl.hpp"
#include "exprs/function/cast/cast_to_string.h"
+#include "storage/olap_common.h"
namespace doris {
@@ -43,6 +46,35 @@ static constexpr int32_t date_threshold = 719528;
namespace {
+constexpr int32_t DORIS_DATE_EPOCH_DAYNR = 719528;
+
+Status decode_date_orc_values(const DataTypeSerDe& serde, IColumn& column,
+ const OrcDecodedColumnView& orc_view) {
+ const auto* orc_batch = dynamic_cast<const
::orc::LongVectorBatch*>(orc_view.batch);
+ if (orc_batch == nullptr) {
+ return Status::InternalError("Unexpected ORC date batch type {}",
+ orc_view.batch->toString());
+ }
+ auto view = orc_serde_utils::make_orc_decoded_view(orc_view,
DecodedValueKind::INT32);
+ NullMap null_map;
+ orc_serde_utils::fill_orc_decoded_null_map(*orc_view.batch, orc_view.rows,
+ orc_view.selected_rows,
&null_map);
+ view.null_map = null_map.empty() ? nullptr : null_map.data();
+ const auto output_rows =
+ orc_serde_utils::orc_decode_row_count(orc_view.rows,
orc_view.selected_rows);
+ std::vector<int32_t> date_values;
+ date_values.resize(output_rows);
+ auto& date_dict = date_day_offset_dict::get();
+ for (size_t row = 0; row < output_rows; ++row) {
+ const auto source_row = orc_serde_utils::orc_source_row_at(row,
orc_view.selected_rows);
+ const auto date =
date_dict[cast_set<int>(orc_batch->data[source_row])];
+ date_values[row] = cast_set<int32_t>(date.daynr() -
DORIS_DATE_EPOCH_DAYNR);
+ }
+ view.values = reinterpret_cast<const uint8_t*>(date_values.data());
+ RETURN_IF_ERROR(orc_serde_utils::read_decoded_values(serde, column,
&view));
+ return Status::OK();
+}
+
Status decode_parquet_date(int32_t encoded_date, DateV2Value<DateV2ValueType>*
value) {
DORIS_CHECK(value != nullptr);
const int64_t day_number = static_cast<int64_t>(encoded_date) +
date_threshold;
@@ -696,4 +728,15 @@ template Status
DataTypeDateV2SerDe::from_decimal_strict_mode_batch<DataTypeDeci
template Status
DataTypeDateV2SerDe::from_decimal_strict_mode_batch<DataTypeDecimal256>(
const DataTypeDecimal256::ColumnType& decimal_col, IColumn&
target_col) const;
+Status DataTypeDateV2SerDe::read_column_from_orc(IColumn& column,
+ const OrcDecodedColumnView&
view) const {
+ DORIS_CHECK(view.file_type != nullptr);
+ DORIS_CHECK(view.batch != nullptr);
+ DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::DATE);
+ if (orc_serde_utils::orc_decode_row_count(view.rows, view.selected_rows)
== 0) {
+ return Status::OK();
+ }
+ return decode_date_orc_values(*this, column, view);
+}
+
} // namespace doris
diff --git a/be/src/core/data_type_serde/data_type_decimal_serde.cpp
b/be/src/core/data_type_serde/data_type_decimal_serde.cpp
index 0b4452e3490..99d32e0f1b8 100644
--- a/be/src/core/data_type_serde/data_type_decimal_serde.cpp
+++ b/be/src/core/data_type_serde/data_type_decimal_serde.cpp
@@ -37,6 +37,7 @@
#include "core/data_type/storage_field_type.h"
#include "core/data_type_serde/arrow_validation.h"
#include "core/data_type_serde/decoded_column_view.h"
+#include "core/data_type_serde/orc_serde_utils.h"
#include "core/data_type_serde/parquet_decode_source.h"
#include "core/types.h"
#include "exec/common/arithmetic_overflow.h"
@@ -52,6 +53,115 @@
namespace doris {
namespace {
+constexpr int DECIMAL_PRECISION_FOR_HIVE11 =
BeConsts::MAX_DECIMAL128_PRECISION;
+
+Int128 to_int128(::orc::Int128 value) {
+ const auto high_bits =
static_cast<__uint128_t>(static_cast<uint64_t>(value.getHighBits()));
+ const auto low_bits = static_cast<__uint128_t>(value.getLowBits());
+ return static_cast<Int128>((high_bits << 64) | low_bits);
+}
+
+::orc::Int128 to_orc_int128(Int128 value) {
+ const auto unsigned_value = static_cast<__uint128_t>(value);
+ return
::orc::Int128(static_cast<int64_t>(static_cast<uint64_t>(unsigned_value >> 64)),
+ static_cast<uint64_t>(unsigned_value));
+}
+
+Status scale_decimal_value(Int128 value, int32_t source_scale, int32_t
target_scale,
+ Int128* scaled_value) {
+ DORIS_CHECK(scaled_value != nullptr);
+ if (source_scale == target_scale) {
+ *scaled_value = value;
+ return Status::OK();
+ }
+ if (source_scale < target_scale) {
+ bool overflow = false;
+ const auto scaled =
::orc::scaleUpInt128ByPowerOfTen(to_orc_int128(value),
+ target_scale -
source_scale, overflow);
+ if (overflow) {
+ return Status::DataQualityError(
+ "ORC decimal value overflows when scaling from {} to {}",
source_scale,
+ target_scale);
+ }
+ *scaled_value = to_int128(scaled);
+ return Status::OK();
+ }
+ *scaled_value = to_int128(
+ ::orc::scaleDownInt128ByPowerOfTen(to_orc_int128(value),
source_scale - target_scale));
+ return Status::OK();
+}
+
+void fill_decimal_big_endian_value(Int128 value, std::array<uint8_t,
sizeof(Int128)>* bytes) {
+ DORIS_CHECK(bytes != nullptr);
+ const auto unsigned_value = static_cast<__uint128_t>(value);
+ for (size_t byte_idx = 0; byte_idx < bytes->size(); ++byte_idx) {
+ const auto shift = (bytes->size() - byte_idx - 1) * 8;
+ (*bytes)[byte_idx] = static_cast<uint8_t>(unsigned_value >> shift);
+ }
+}
+
+Status decode_decimal_orc_values(const DataTypeSerDe& serde, IColumn& column,
+ const OrcDecodedColumnView& orc_view, int32_t
target_scale) {
+ DORIS_CHECK(orc_view.file_type != nullptr);
+ auto view = orc_serde_utils::make_orc_decoded_view(orc_view,
DecodedValueKind::FIXED_BINARY);
+ view.decimal_precision = orc_view.file_type->getPrecision() == 0
+ ? DECIMAL_PRECISION_FOR_HIVE11
+ :
cast_set<int>(orc_view.file_type->getPrecision());
+ NullMap null_map;
+ orc_serde_utils::fill_orc_decoded_null_map(*orc_view.batch, orc_view.rows,
+ orc_view.selected_rows,
&null_map);
+ view.null_map = null_map.empty() ? nullptr : null_map.data();
+ view.fixed_length = sizeof(Int128);
+
+ std::vector<StringRef> binary_values;
+ std::vector<std::array<uint8_t, sizeof(Int128)>> decimal_values;
+ const auto output_rows =
+ orc_serde_utils::orc_decode_row_count(orc_view.rows,
orc_view.selected_rows);
+ decimal_values.resize(output_rows);
+ binary_values.reserve(output_rows);
+ if (const auto* decimal64_batch =
+ dynamic_cast<const
::orc::Decimal64VectorBatch*>(orc_view.batch);
+ decimal64_batch != nullptr) {
+ view.decimal_scale = decimal64_batch->scale;
+ for (size_t row = 0; row < output_rows; ++row) {
+ Int128 value = 0;
+ const auto source_row = orc_serde_utils::orc_source_row_at(row,
orc_view.selected_rows);
+ if (!orc_serde_utils::orc_row_is_null(*orc_view.batch,
source_row)) {
+
RETURN_IF_ERROR(scale_decimal_value(decimal64_batch->values[source_row],
+ decimal64_batch->scale,
target_scale, &value));
+ }
+ fill_decimal_big_endian_value(value, &decimal_values[row]);
+ binary_values.emplace_back(reinterpret_cast<const
char*>(decimal_values[row].data()),
+ decimal_values[row].size());
+ }
+ view.binary_values = &binary_values;
+ RETURN_IF_ERROR(orc_serde_utils::read_decoded_values(serde, column,
&view));
+ return Status::OK();
+ }
+
+ const auto* decimal128_batch =
+ dynamic_cast<const ::orc::Decimal128VectorBatch*>(orc_view.batch);
+ if (decimal128_batch == nullptr) {
+ return Status::InternalError("Unexpected ORC decimal batch type {}",
+ orc_view.batch->toString());
+ }
+ view.decimal_scale = decimal128_batch->scale;
+ for (size_t row = 0; row < output_rows; ++row) {
+ Int128 value = 0;
+ const auto source_row = orc_serde_utils::orc_source_row_at(row,
orc_view.selected_rows);
+ if (!orc_serde_utils::orc_row_is_null(*orc_view.batch, source_row)) {
+
RETURN_IF_ERROR(scale_decimal_value(to_int128(decimal128_batch->values[source_row]),
+ decimal128_batch->scale,
target_scale, &value));
+ }
+ fill_decimal_big_endian_value(value, &decimal_values[row]);
+ binary_values.emplace_back(reinterpret_cast<const
char*>(decimal_values[row].data()),
+ decimal_values[row].size());
+ }
+ view.binary_values = &binary_values;
+ RETURN_IF_ERROR(orc_serde_utils::read_decoded_values(serde, column,
&view));
+ return Status::OK();
+}
+
template <typename NativeType>
NativeType decode_big_endian_signed_integer(const uint8_t* data, int length) {
if constexpr (std::is_same_v<NativeType, wide::Int256>) {
@@ -1360,6 +1470,18 @@ const uint8_t*
DataTypeDecimalSerDe<T>::deserialize_binary_to_field(const uint8_
return data;
}
+template <PrimitiveType T>
+Status DataTypeDecimalSerDe<T>::read_column_from_orc(IColumn& column,
+ const
OrcDecodedColumnView& view) const {
+ DORIS_CHECK(view.file_type != nullptr);
+ DORIS_CHECK(view.batch != nullptr);
+ DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::DECIMAL);
+ if (orc_serde_utils::orc_decode_row_count(view.rows, view.selected_rows)
== 0) {
+ return Status::OK();
+ }
+ return decode_decimal_orc_values(*this, column, view,
cast_set<int32_t>(scale));
+}
+
template class DataTypeDecimalSerDe<TYPE_DECIMAL32>;
template class DataTypeDecimalSerDe<TYPE_DECIMAL64>;
template class DataTypeDecimalSerDe<TYPE_DECIMAL128I>;
diff --git a/be/src/core/data_type_serde/data_type_map_serde.cpp
b/be/src/core/data_type_serde/data_type_map_serde.cpp
index cd1ec5cf984..419006257a5 100644
--- a/be/src/core/data_type_serde/data_type_map_serde.cpp
+++ b/be/src/core/data_type_serde/data_type_map_serde.cpp
@@ -36,6 +36,59 @@
namespace doris {
class Arena;
+
+namespace {
+
+Status decode_map_orc_values(const DataTypeSerDeSPtr& key_serde,
+ const DataTypeSerDeSPtr& value_serde, IColumn&
nested_column,
+ const OrcDecodedColumnView& orc_view) {
+ const auto* orc_map = dynamic_cast<const
::orc::MapVectorBatch*>(orc_view.batch);
+ if (orc_map == nullptr) {
+ return Status::InternalError("Unexpected ORC map batch type {}",
+ orc_view.batch->toString());
+ }
+ DORIS_CHECK(orc_view.file_type != nullptr);
+ DORIS_CHECK(orc_view.selected_type != nullptr);
+ DORIS_CHECK(orc_view.file_type->getSubtypeCount() == 2);
+ DORIS_CHECK(orc_view.selected_type->getSubtypeCount() == 2);
+ DORIS_CHECK(orc_map->keys != nullptr);
+ DORIS_CHECK(orc_map->elements != nullptr);
+ auto& map_column = assert_cast<ColumnMap&>(nested_column);
+ size_t element_size = 0;
+ std::vector<size_t> element_selection;
+ RETURN_IF_ERROR(orc_serde_utils::append_orc_offsets(
+ map_column.get_offsets(), orc_map->offsets, orc_view.rows,
&element_size,
+ orc_view.selected_rows, &element_selection));
+ const auto* file_key_type = orc_view.file_type->getSubtype(0);
+ const auto* selected_key_type = orc_view.selected_type->getSubtype(0);
+ DORIS_CHECK(file_key_type != nullptr);
+ DORIS_CHECK(selected_key_type != nullptr);
+ const auto child_rows = orc_view.selected_rows == nullptr
+ ? element_size
+ :
static_cast<size_t>(orc_map->keys->numElements);
+ const auto* child_selection = orc_view.selected_rows == nullptr ? nullptr
: &element_selection;
+ auto key_column = map_column.get_keys_ptr()->assert_mutable();
+ auto key_view =
+ orc_serde_utils::make_child_orc_view(orc_view, file_key_type,
selected_key_type,
+ orc_map->keys.get(),
child_rows, child_selection);
+ RETURN_IF_ERROR(orc_serde_utils::read_orc_child_column(key_serde,
key_column, key_view));
+ map_column.get_keys_ptr() = std::move(key_column);
+ const auto* file_value_type = orc_view.file_type->getSubtype(1);
+ const auto* selected_value_type = orc_view.selected_type->getSubtype(1);
+ DORIS_CHECK(file_value_type != nullptr);
+ DORIS_CHECK(selected_value_type != nullptr);
+ auto value_column = map_column.get_values_ptr()->assert_mutable();
+ auto value_view = orc_serde_utils::make_child_orc_view(
+ orc_view, file_value_type, selected_value_type,
orc_map->elements.get(),
+ orc_view.selected_rows == nullptr ? element_size
+ :
static_cast<size_t>(orc_map->elements->numElements),
+ child_selection);
+ RETURN_IF_ERROR(orc_serde_utils::read_orc_child_column(value_serde,
value_column, value_view));
+ map_column.get_values_ptr() = std::move(value_column);
+ return Status::OK();
+}
+
+} // namespace
Status DataTypeMapSerDe::serialize_column_to_json(const IColumn& column,
int64_t start_idx,
int64_t end_idx,
BufferWritable& bw,
FormatOptions& options)
const {
@@ -484,8 +537,8 @@ Status DataTypeMapSerDe::write_column_to_orc(const
std::string& timezone, const
packed_nested_size,
arena, options));
// String batches borrow their source bytes, but the packed columns are
local to this call;
// keep only those borrowed leaves in the write Arena until Writer::add()
consumes them.
- copy_orc_string_data_to_arena(cur_batch->keys.get(), arena);
- copy_orc_string_data_to_arena(cur_batch->elements.get(), arena);
+ orc_serde_utils::copy_orc_string_data_to_arena(cur_batch->keys.get(),
arena);
+ orc_serde_utils::copy_orc_string_data_to_arena(cur_batch->elements.get(),
arena);
cur_batch->keys->numElements = packed_nested_size;
cur_batch->elements->numElements = packed_nested_size;
@@ -722,4 +775,15 @@ bool DataTypeMapSerDe::write_column_to_hive_text(const
IColumn& column, BufferWr
return true;
}
+Status DataTypeMapSerDe::read_column_from_orc(IColumn& column,
+ const OrcDecodedColumnView&
view) const {
+ DORIS_CHECK(view.file_type != nullptr);
+ DORIS_CHECK(view.batch != nullptr);
+ DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::MAP);
+ if (orc_serde_utils::orc_decode_row_count(view.rows, view.selected_rows)
== 0) {
+ return Status::OK();
+ }
+ return decode_map_orc_values(key_serde, value_serde, column, view);
+}
+
} // namespace doris
diff --git a/be/src/core/data_type_serde/data_type_nullable_serde.cpp
b/be/src/core/data_type_serde/data_type_nullable_serde.cpp
index 2b4b6870581..c06bcf0af82 100644
--- a/be/src/core/data_type_serde/data_type_nullable_serde.cpp
+++ b/be/src/core/data_type_serde/data_type_nullable_serde.cpp
@@ -22,6 +22,7 @@
#include <algorithm>
#include <boost/iterator/iterator_facade.hpp>
+#include <cstring>
#include <vector>
#include "common/config.h"
@@ -34,6 +35,7 @@
#include "core/data_type_serde/data_type_serde.h"
#include "core/data_type_serde/data_type_string_serde.h"
#include "core/data_type_serde/decoded_column_view.h"
+#include "core/data_type_serde/orc_serde_utils.h"
#include "exprs/function/cast/cast_base.h"
#include "format/transformer/vcsv_transformer.h"
#include "util/jsonb_document.h"
@@ -42,6 +44,26 @@
namespace doris {
class Arena;
+
+namespace {
+
+void append_orc_null_map(const ::orc::ColumnVectorBatch& batch, size_t rows,
+ const std::vector<size_t>* selected_rows, NullMap*
null_map) {
+ DORIS_CHECK(null_map != nullptr);
+ const auto output_rows = orc_serde_utils::orc_decode_row_count(rows,
selected_rows);
+ const auto old_size = null_map->size();
+ null_map->resize(old_size + output_rows);
+ if (batch.hasNulls) {
+ for (size_t row = 0; row < output_rows; ++row) {
+ (*null_map)[old_size + row] =
+ !batch.notNull[orc_serde_utils::orc_source_row_at(row,
selected_rows)];
+ }
+ return;
+ }
+ std::memset(null_map->data() + old_size, 0, output_rows);
+}
+
+} // namespace
Status DataTypeNullableSerDe::serialize_column_to_json(const IColumn& column,
int64_t start_idx,
int64_t end_idx,
BufferWritable& bw,
FormatOptions& options)
const {
@@ -591,4 +613,29 @@ Status
DataTypeNullableSerDe::from_string_strict_mode(StringRef& str, IColumn& c
null_column.get_null_map_data().push_back(0);
return Status::OK();
}
+
+Status DataTypeNullableSerDe::read_column_from_orc(IColumn& column,
+ const OrcDecodedColumnView&
view) const {
+ DORIS_CHECK(view.file_type != nullptr);
+ DORIS_CHECK(view.selected_type != nullptr);
+ DORIS_CHECK(view.batch != nullptr);
+ DORIS_CHECK(view.file_type->getKind() == view.selected_type->getKind());
+ auto& nullable_column = assert_cast<ColumnNullable&>(column);
+ const auto output_rows = orc_serde_utils::orc_decode_row_count(view.rows,
view.selected_rows);
+ if (output_rows == 0) {
+ return Status::OK();
+ }
+
+ auto& null_map = nullable_column.get_null_map_data();
+ const auto old_null_map_size = null_map.size();
+ auto& nested_column = nullable_column.get_nested_column();
+ const auto old_nested_size = nested_column.size();
+ append_orc_null_map(*view.batch, view.rows, view.selected_rows, &null_map);
+ auto st = nested_serde->read_column_from_orc(nested_column, view);
+ if (!st.ok()) {
+ null_map.resize(old_null_map_size);
+ nested_column.resize(old_nested_size);
+ }
+ return st;
+}
} // namespace doris
diff --git a/be/src/core/data_type_serde/data_type_number_serde.cpp
b/be/src/core/data_type_serde/data_type_number_serde.cpp
index df9ccb529da..ca30ec96d25 100644
--- a/be/src/core/data_type_serde/data_type_number_serde.cpp
+++ b/be/src/core/data_type_serde/data_type_number_serde.cpp
@@ -23,7 +23,9 @@
#include <cmath>
#include <cstdint>
#include <limits>
+#include <memory>
#include <type_traits>
+#include <vector>
#include "common/config.h"
#include "common/exception.h"
@@ -35,6 +37,7 @@
#include "core/data_type_serde/arrow_validation.h"
#include "core/data_type_serde/data_type_serde.h"
#include "core/data_type_serde/decoded_column_view.h"
+#include "core/data_type_serde/orc_serde_utils.h"
#include "core/data_type_serde/parquet_decode_source.h"
#include "core/packed_int128.h"
#include "core/types.h"
@@ -54,6 +57,95 @@
namespace doris {
namespace {
+template <typename SourceType>
+void fill_selected_values(const SourceType* source_values, size_t rows,
+ const std::vector<size_t>* selected_rows,
+ std::vector<SourceType>* selected_values) {
+ DORIS_CHECK(source_values != nullptr);
+ DORIS_CHECK(selected_values != nullptr);
+ const auto output_rows = orc_serde_utils::orc_decode_row_count(rows,
selected_rows);
+ selected_values->resize(output_rows);
+ for (size_t row = 0; row < output_rows; ++row) {
+ (*selected_values)[row] =
+ source_values[orc_serde_utils::orc_source_row_at(row,
selected_rows)];
+ }
+}
+
+template <typename OrcBatchType, typename SourceType>
+Status decode_fixed_orc_values(const DataTypeSerDe& serde, IColumn& column,
+ const OrcDecodedColumnView& orc_view,
DecodedValueKind value_kind) {
+ const auto* orc_batch = dynamic_cast<const OrcBatchType*>(orc_view.batch);
+ if (orc_batch == nullptr) {
+ return Status::InternalError("Unexpected ORC scalar batch type {}",
+ orc_view.batch->toString());
+ }
+ auto view = orc_serde_utils::make_orc_decoded_view(orc_view, value_kind);
+ NullMap null_map;
+ orc_serde_utils::fill_orc_decoded_null_map(*orc_view.batch, orc_view.rows,
+ orc_view.selected_rows,
&null_map);
+ view.null_map = null_map.empty() ? nullptr : null_map.data();
+ std::vector<SourceType> selected_values;
+ if (orc_view.selected_rows == nullptr) {
+ view.values = reinterpret_cast<const uint8_t*>(orc_batch->data.data());
+ } else {
+ fill_selected_values(orc_batch->data.data(), orc_view.rows,
orc_view.selected_rows,
+ &selected_values);
+ view.values = reinterpret_cast<const uint8_t*>(selected_values.data());
+ }
+ RETURN_IF_ERROR(orc_serde_utils::read_decoded_values(serde, column,
&view));
+ return Status::OK();
+}
+
+Status decode_float_orc_values(const DataTypeSerDe& serde, IColumn& column,
+ const OrcDecodedColumnView& orc_view) {
+ const auto* orc_batch = dynamic_cast<const
::orc::DoubleVectorBatch*>(orc_view.batch);
+ if (orc_batch == nullptr) {
+ return Status::InternalError("Unexpected ORC float batch type {}",
+ orc_view.batch->toString());
+ }
+ auto view = orc_serde_utils::make_orc_decoded_view(orc_view,
DecodedValueKind::FLOAT);
+ NullMap null_map;
+ orc_serde_utils::fill_orc_decoded_null_map(*orc_view.batch, orc_view.rows,
+ orc_view.selected_rows,
&null_map);
+ view.null_map = null_map.empty() ? nullptr : null_map.data();
+ const auto output_rows =
+ orc_serde_utils::orc_decode_row_count(orc_view.rows,
orc_view.selected_rows);
+ std::vector<float> float_values;
+ float_values.resize(output_rows);
+ for (size_t row = 0; row < output_rows; ++row) {
+ float_values[row] = static_cast<float>(
+ orc_batch->data[orc_serde_utils::orc_source_row_at(row,
orc_view.selected_rows)]);
+ }
+ view.values = reinterpret_cast<const uint8_t*>(float_values.data());
+ RETURN_IF_ERROR(orc_serde_utils::read_decoded_values(serde, column,
&view));
+ return Status::OK();
+}
+
+Status decode_boolean_orc_values(const DataTypeSerDe& serde, IColumn& column,
+ const OrcDecodedColumnView& orc_view) {
+ const auto* orc_batch = dynamic_cast<const
::orc::LongVectorBatch*>(orc_view.batch);
+ if (orc_batch == nullptr) {
+ return Status::InternalError("Unexpected ORC boolean batch type {}",
+ orc_view.batch->toString());
+ }
+ auto view = orc_serde_utils::make_orc_decoded_view(orc_view,
DecodedValueKind::BOOL);
+ NullMap null_map;
+ orc_serde_utils::fill_orc_decoded_null_map(*orc_view.batch, orc_view.rows,
+ orc_view.selected_rows,
&null_map);
+ view.null_map = null_map.empty() ? nullptr : null_map.data();
+ const auto output_rows =
+ orc_serde_utils::orc_decode_row_count(orc_view.rows,
orc_view.selected_rows);
+ std::unique_ptr<bool[]> bool_values =
std::make_unique<bool[]>(output_rows);
+ for (size_t row = 0; row < output_rows; ++row) {
+ bool_values[row] =
+ orc_batch->data[orc_serde_utils::orc_source_row_at(row,
orc_view.selected_rows)] !=
+ 0;
+ }
+ view.values = reinterpret_cast<const uint8_t*>(bool_values.get());
+ RETURN_IF_ERROR(orc_serde_utils::read_decoded_values(serde, column,
&view));
+ return Status::OK();
+}
+
float parquet_half_to_float(uint16_t half) {
const uint32_t sign = (half & 0x8000U) << 16;
const uint32_t exponent = (half & 0x7C00U) >> 10;
@@ -1747,6 +1839,42 @@ void DataTypeNumberSerDe<T>::to_string_batch(const
IColumn& column, ColumnString
}
}
+template <PrimitiveType T>
+Status DataTypeNumberSerDe<T>::read_column_from_orc(IColumn& column,
+ const
OrcDecodedColumnView& view) const {
+ DORIS_CHECK(view.file_type != nullptr);
+ DORIS_CHECK(view.batch != nullptr);
+ if (orc_serde_utils::orc_decode_row_count(view.rows, view.selected_rows)
== 0) {
+ return Status::OK();
+ }
+
+ if constexpr (T == TYPE_BOOLEAN) {
+ DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::BOOLEAN);
+ return decode_boolean_orc_values(*this, column, view);
+ } else if constexpr (T == TYPE_TINYINT || T == TYPE_SMALLINT || T ==
TYPE_INT ||
+ T == TYPE_BIGINT) {
+ if constexpr (T == TYPE_TINYINT) {
+ DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::BYTE);
+ } else if constexpr (T == TYPE_SMALLINT) {
+ DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::SHORT);
+ } else if constexpr (T == TYPE_INT) {
+ DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::INT);
+ } else {
+ DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::LONG);
+ }
+ return decode_fixed_orc_values<::orc::LongVectorBatch, int64_t>(*this,
column, view,
+
DecodedValueKind::INT64);
+ } else if constexpr (T == TYPE_FLOAT) {
+ DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::FLOAT);
+ return decode_float_orc_values(*this, column, view);
+ } else if constexpr (T == TYPE_DOUBLE) {
+ DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::DOUBLE);
+ return decode_fixed_orc_values<::orc::DoubleVectorBatch,
double>(*this, column, view,
+
DecodedValueKind::DOUBLE);
+ }
+ return DataTypeSerDe::read_column_from_orc(column, view);
+}
+
/// Explicit template instantiations - to avoid code bloat in headers.
template class DataTypeNumberSerDe<TYPE_BOOLEAN>;
template class DataTypeNumberSerDe<TYPE_TINYINT>;
diff --git a/be/src/core/data_type_serde/data_type_serde.cpp
b/be/src/core/data_type_serde/data_type_serde.cpp
index ba5426681c4..ca5f2c206ad 100644
--- a/be/src/core/data_type_serde/data_type_serde.cpp
+++ b/be/src/core/data_type_serde/data_type_serde.cpp
@@ -16,722 +16,28 @@
// under the License.
#include "core/data_type_serde/data_type_serde.h"
-#include <cctz/time_zone.h>
-
-#include <array>
-#include <cstring>
-#include <memory>
-#include <orc/OrcFile.hh>
-#include <orc/Vector.hh>
#include <vector>
-#include "common/cast_set.h"
#include "common/check.h"
-#include "common/consts.h"
#include "common/exception.h"
#include "common/status.h"
#include "core/assert_cast.h"
#include "core/column/column.h"
-#include "core/column/column_array.h"
-#include "core/column/column_map.h"
#include "core/column/column_nullable.h"
-#include "core/column/column_struct.h"
-#include "core/column/column_vector.h"
#include "core/data_type/data_type.h"
-#include "core/data_type/data_type_array.h"
-#include "core/data_type/data_type_map.h"
-#include "core/data_type/data_type_nullable.h"
-#include "core/data_type/data_type_struct.h"
#include "core/data_type/storage_field_type.h"
#include "core/data_type_serde/data_type_array_serde.h"
-#include "core/data_type_serde/data_type_datetimev2_serde.h"
-#include "core/data_type_serde/data_type_datev2_serde.h"
#include "core/data_type_serde/data_type_decimal_serde.h"
#include "core/data_type_serde/data_type_jsonb_serde.h"
-#include "core/data_type_serde/data_type_nullable_serde.h"
#include "core/data_type_serde/data_type_number_serde.h"
#include "core/data_type_serde/data_type_string_serde.h"
-#include "core/data_type_serde/data_type_timestamptz_serde.h"
#include "core/data_type_serde/parquet_decode_source.h"
#include "core/field.h"
-#include "core/types.h"
-#include "core/value/timestamptz_value.h"
-#include "core/value/vdatetime_value.h"
#include "exprs/function/cast/cast_base.h"
#include "runtime/descriptors.h"
-#include "storage/olap_common.h"
#include "util/jsonb_document.h"
#include "util/jsonb_writer.h"
namespace doris {
-namespace {
-
-constexpr int DECIMAL_PRECISION_FOR_HIVE11 =
BeConsts::MAX_DECIMAL128_PRECISION;
-constexpr int32_t DORIS_DATE_EPOCH_DAYNR = 719528;
-
-bool orc_row_is_null(const ::orc::ColumnVectorBatch& batch, size_t row) {
- return batch.hasNulls && !batch.notNull[row];
-}
-
-size_t orc_decode_row_count(size_t rows, const std::vector<size_t>*
selected_rows) {
- if (selected_rows == nullptr) {
- return rows;
- }
- return selected_rows->size();
-}
-
-size_t orc_source_row_at(size_t row, const std::vector<size_t>* selected_rows)
{
- if (selected_rows == nullptr) {
- return row;
- }
- return (*selected_rows)[row];
-}
-
-DecodedColumnView make_orc_decoded_view(const OrcDecodedColumnView& orc_view,
- DecodedValueKind value_kind) {
- DecodedColumnView view;
- view.value_kind = value_kind;
- view.row_count = cast_set<int64_t>(orc_decode_row_count(orc_view.rows,
orc_view.selected_rows));
- view.timezone = orc_view.timezone;
- return view;
-}
-
-void fill_orc_decoded_null_map(const ::orc::ColumnVectorBatch& batch, size_t
rows,
- const std::vector<size_t>* selected_rows,
NullMap* null_map) {
- DORIS_CHECK(null_map != nullptr);
- if (!batch.hasNulls) {
- return;
- }
- const auto output_rows = orc_decode_row_count(rows, selected_rows);
- null_map->resize(output_rows);
- for (size_t row = 0; row < output_rows; ++row) {
- (*null_map)[row] = !batch.notNull[orc_source_row_at(row,
selected_rows)];
- }
-}
-
-void append_orc_null_map(const ::orc::ColumnVectorBatch& batch, size_t rows,
- const std::vector<size_t>* selected_rows, NullMap*
null_map) {
- DORIS_CHECK(null_map != nullptr);
- const auto output_rows = orc_decode_row_count(rows, selected_rows);
- const auto old_size = null_map->size();
- null_map->resize(old_size + output_rows);
- if (batch.hasNulls) {
- for (size_t row = 0; row < output_rows; ++row) {
- (*null_map)[old_size + row] =
!batch.notNull[orc_source_row_at(row, selected_rows)];
- }
- return;
- }
- std::memset(null_map->data() + old_size, 0, output_rows);
-}
-
-size_t trim_right_spaces(const char* value, size_t length) {
- while (length > 0 && value[length - 1] == ' ') {
- --length;
- }
- return length;
-}
-
-Status append_orc_string_ref(const ::orc::Type& file_type, const char* data,
int64_t length,
- std::vector<StringRef>& binary_values) {
- if (length < 0) {
- return Status::Corruption("Invalid negative ORC string length {}",
length);
- }
- auto value_length = static_cast<size_t>(length);
- if (file_type.getKind() == ::orc::TypeKind::CHAR) {
- value_length = trim_right_spaces(data, value_length);
- }
- binary_values.emplace_back(value_length == 0 ? "" : data, value_length);
- return Status::OK();
-}
-
-Int128 to_int128(::orc::Int128 value) {
- const auto high_bits =
static_cast<__uint128_t>(static_cast<uint64_t>(value.getHighBits()));
- const auto low_bits = static_cast<__uint128_t>(value.getLowBits());
- return static_cast<Int128>((high_bits << 64) | low_bits);
-}
-
-::orc::Int128 to_orc_int128(Int128 value) {
- const auto unsigned_value = static_cast<__uint128_t>(value);
- return
::orc::Int128(static_cast<int64_t>(static_cast<uint64_t>(unsigned_value >> 64)),
- static_cast<uint64_t>(unsigned_value));
-}
-
-Status scale_decimal_value(Int128 value, int32_t source_scale, int32_t
target_scale,
- Int128* scaled_value) {
- DORIS_CHECK(scaled_value != nullptr);
- if (source_scale == target_scale) {
- *scaled_value = value;
- return Status::OK();
- }
- if (source_scale < target_scale) {
- bool overflow = false;
- const auto scaled =
::orc::scaleUpInt128ByPowerOfTen(to_orc_int128(value),
- target_scale -
source_scale, overflow);
- if (overflow) {
- return Status::DataQualityError(
- "ORC decimal value overflows when scaling from {} to {}",
source_scale,
- target_scale);
- }
- *scaled_value = to_int128(scaled);
- return Status::OK();
- }
- *scaled_value = to_int128(
- ::orc::scaleDownInt128ByPowerOfTen(to_orc_int128(value),
source_scale - target_scale));
- return Status::OK();
-}
-
-void fill_decimal_big_endian_value(Int128 value, std::array<uint8_t,
sizeof(Int128)>* bytes) {
- DORIS_CHECK(bytes != nullptr);
- const auto unsigned_value = static_cast<__uint128_t>(value);
- for (size_t byte_idx = 0; byte_idx < bytes->size(); ++byte_idx) {
- const auto shift = (bytes->size() - byte_idx - 1) * 8;
- (*bytes)[byte_idx] = static_cast<uint8_t>(unsigned_value >> shift);
- }
-}
-
-Status read_decoded_values(const DataTypeSerDe& serde, IColumn& column,
DecodedColumnView* view) {
- DORIS_CHECK(view != nullptr);
- RETURN_IF_ERROR(serde.read_column_from_decoded_values(column, *view));
- return Status::OK();
-}
-
-template <typename SourceType>
-void fill_selected_values(const SourceType* source_values, size_t rows,
- const std::vector<size_t>* selected_rows,
- std::vector<SourceType>* selected_values) {
- DORIS_CHECK(source_values != nullptr);
- DORIS_CHECK(selected_values != nullptr);
- const auto output_rows = orc_decode_row_count(rows, selected_rows);
- selected_values->resize(output_rows);
- for (size_t row = 0; row < output_rows; ++row) {
- (*selected_values)[row] = source_values[orc_source_row_at(row,
selected_rows)];
- }
-}
-
-template <typename OrcBatchType, typename SourceType>
-Status decode_fixed_orc_values(const DataTypeSerDe& serde, IColumn& column,
- const OrcDecodedColumnView& orc_view,
DecodedValueKind value_kind) {
- const auto* orc_batch = dynamic_cast<const OrcBatchType*>(orc_view.batch);
- if (orc_batch == nullptr) {
- return Status::InternalError("Unexpected ORC scalar batch type {}",
- orc_view.batch->toString());
- }
- auto view = make_orc_decoded_view(orc_view, value_kind);
- NullMap null_map;
- fill_orc_decoded_null_map(*orc_view.batch, orc_view.rows,
orc_view.selected_rows, &null_map);
- view.null_map = null_map.empty() ? nullptr : null_map.data();
- std::vector<SourceType> selected_values;
- if (orc_view.selected_rows == nullptr) {
- view.values = reinterpret_cast<const uint8_t*>(orc_batch->data.data());
- } else {
- fill_selected_values(orc_batch->data.data(), orc_view.rows,
orc_view.selected_rows,
- &selected_values);
- view.values = reinterpret_cast<const uint8_t*>(selected_values.data());
- }
- RETURN_IF_ERROR(read_decoded_values(serde, column, &view));
- return Status::OK();
-}
-
-Status decode_float_orc_values(const DataTypeSerDe& serde, IColumn& column,
- const OrcDecodedColumnView& orc_view) {
- const auto* orc_batch = dynamic_cast<const
::orc::DoubleVectorBatch*>(orc_view.batch);
- if (orc_batch == nullptr) {
- return Status::InternalError("Unexpected ORC float batch type {}",
- orc_view.batch->toString());
- }
- auto view = make_orc_decoded_view(orc_view, DecodedValueKind::FLOAT);
- NullMap null_map;
- fill_orc_decoded_null_map(*orc_view.batch, orc_view.rows,
orc_view.selected_rows, &null_map);
- view.null_map = null_map.empty() ? nullptr : null_map.data();
- const auto output_rows = orc_decode_row_count(orc_view.rows,
orc_view.selected_rows);
- std::vector<float> float_values;
- float_values.resize(output_rows);
- for (size_t row = 0; row < output_rows; ++row) {
- float_values[row] =
- static_cast<float>(orc_batch->data[orc_source_row_at(row,
orc_view.selected_rows)]);
- }
- view.values = reinterpret_cast<const uint8_t*>(float_values.data());
- RETURN_IF_ERROR(read_decoded_values(serde, column, &view));
- return Status::OK();
-}
-
-Status decode_boolean_orc_values(const DataTypeSerDe& serde, IColumn& column,
- const OrcDecodedColumnView& orc_view) {
- const auto* orc_batch = dynamic_cast<const
::orc::LongVectorBatch*>(orc_view.batch);
- if (orc_batch == nullptr) {
- return Status::InternalError("Unexpected ORC boolean batch type {}",
- orc_view.batch->toString());
- }
- auto view = make_orc_decoded_view(orc_view, DecodedValueKind::BOOL);
- NullMap null_map;
- fill_orc_decoded_null_map(*orc_view.batch, orc_view.rows,
orc_view.selected_rows, &null_map);
- view.null_map = null_map.empty() ? nullptr : null_map.data();
- const auto output_rows = orc_decode_row_count(orc_view.rows,
orc_view.selected_rows);
- std::unique_ptr<bool[]> bool_values =
std::make_unique<bool[]>(output_rows);
- for (size_t row = 0; row < output_rows; ++row) {
- bool_values[row] = orc_batch->data[orc_source_row_at(row,
orc_view.selected_rows)] != 0;
- }
- view.values = reinterpret_cast<const uint8_t*>(bool_values.get());
- RETURN_IF_ERROR(read_decoded_values(serde, column, &view));
- return Status::OK();
-}
-
-Status decode_string_orc_values(const DataTypeSerDe& serde, IColumn& column,
- const OrcDecodedColumnView& orc_view) {
- DORIS_CHECK(orc_view.file_type != nullptr);
- if (const auto* encoded_batch =
- dynamic_cast<const
::orc::EncodedStringVectorBatch*>(orc_view.batch);
- encoded_batch != nullptr && encoded_batch->isEncoded) {
- if (encoded_batch->dictionary == nullptr) {
- return Status::InternalError("Encoded ORC string batch has no
dictionary");
- }
- auto view = make_orc_decoded_view(orc_view, DecodedValueKind::BINARY);
- NullMap null_map;
- fill_orc_decoded_null_map(*orc_view.batch, orc_view.rows,
orc_view.selected_rows,
- &null_map);
- view.null_map = null_map.empty() ? nullptr : null_map.data();
- const auto output_rows = orc_decode_row_count(orc_view.rows,
orc_view.selected_rows);
- std::vector<StringRef> binary_values;
- binary_values.reserve(output_rows);
- for (size_t row = 0; row < output_rows; ++row) {
- const auto source_row = orc_source_row_at(row,
orc_view.selected_rows);
- if (orc_row_is_null(*orc_view.batch, source_row)) {
- binary_values.emplace_back("", 0);
- continue;
- }
- char* data = nullptr;
- int64_t length = 0;
-
encoded_batch->dictionary->getValueByIndex(encoded_batch->index[source_row],
data,
- length);
- RETURN_IF_ERROR(
- append_orc_string_ref(*orc_view.file_type, data, length,
binary_values));
- }
- view.binary_values = &binary_values;
- RETURN_IF_ERROR(read_decoded_values(serde, column, &view));
- return Status::OK();
- }
-
- const auto* orc_batch = dynamic_cast<const
::orc::StringVectorBatch*>(orc_view.batch);
- if (orc_batch == nullptr) {
- return Status::InternalError("Unexpected ORC string batch type {}",
- orc_view.batch->toString());
- }
- auto view = make_orc_decoded_view(orc_view, DecodedValueKind::BINARY);
- NullMap null_map;
- fill_orc_decoded_null_map(*orc_view.batch, orc_view.rows,
orc_view.selected_rows, &null_map);
- view.null_map = null_map.empty() ? nullptr : null_map.data();
- const auto output_rows = orc_decode_row_count(orc_view.rows,
orc_view.selected_rows);
- std::vector<StringRef> binary_values;
- binary_values.reserve(output_rows);
- for (size_t row = 0; row < output_rows; ++row) {
- const auto source_row = orc_source_row_at(row, orc_view.selected_rows);
- if (orc_row_is_null(*orc_view.batch, source_row)) {
- binary_values.emplace_back("", 0);
- continue;
- }
- RETURN_IF_ERROR(append_orc_string_ref(*orc_view.file_type,
orc_batch->data[source_row],
- orc_batch->length[source_row],
binary_values));
- }
- view.binary_values = &binary_values;
- RETURN_IF_ERROR(read_decoded_values(serde, column, &view));
- return Status::OK();
-}
-
-Status decode_date_orc_values(const DataTypeSerDe& serde, IColumn& column,
- const OrcDecodedColumnView& orc_view) {
- const auto* orc_batch = dynamic_cast<const
::orc::LongVectorBatch*>(orc_view.batch);
- if (orc_batch == nullptr) {
- return Status::InternalError("Unexpected ORC date batch type {}",
- orc_view.batch->toString());
- }
- auto view = make_orc_decoded_view(orc_view, DecodedValueKind::INT32);
- NullMap null_map;
- fill_orc_decoded_null_map(*orc_view.batch, orc_view.rows,
orc_view.selected_rows, &null_map);
- view.null_map = null_map.empty() ? nullptr : null_map.data();
- const auto output_rows = orc_decode_row_count(orc_view.rows,
orc_view.selected_rows);
- std::vector<int32_t> date_values;
- date_values.resize(output_rows);
- auto& date_dict = date_day_offset_dict::get();
- for (size_t row = 0; row < output_rows; ++row) {
- const auto source_row = orc_source_row_at(row, orc_view.selected_rows);
- const auto date =
date_dict[cast_set<int>(orc_batch->data[source_row])];
- date_values[row] = cast_set<int32_t>(date.daynr() -
DORIS_DATE_EPOCH_DAYNR);
- }
- view.values = reinterpret_cast<const uint8_t*>(date_values.data());
- RETURN_IF_ERROR(read_decoded_values(serde, column, &view));
- return Status::OK();
-}
-
-Status decode_decimal_orc_values(const DataTypeSerDe& serde, IColumn& column,
- const OrcDecodedColumnView& orc_view, int32_t
target_scale) {
- DORIS_CHECK(orc_view.file_type != nullptr);
- auto view = make_orc_decoded_view(orc_view,
DecodedValueKind::FIXED_BINARY);
- view.decimal_precision = orc_view.file_type->getPrecision() == 0
- ? DECIMAL_PRECISION_FOR_HIVE11
- :
cast_set<int>(orc_view.file_type->getPrecision());
- NullMap null_map;
- fill_orc_decoded_null_map(*orc_view.batch, orc_view.rows,
orc_view.selected_rows, &null_map);
- view.null_map = null_map.empty() ? nullptr : null_map.data();
- view.fixed_length = sizeof(Int128);
-
- std::vector<StringRef> binary_values;
- std::vector<std::array<uint8_t, sizeof(Int128)>> decimal_values;
- const auto output_rows = orc_decode_row_count(orc_view.rows,
orc_view.selected_rows);
- decimal_values.resize(output_rows);
- binary_values.reserve(output_rows);
- if (const auto* decimal64_batch =
- dynamic_cast<const
::orc::Decimal64VectorBatch*>(orc_view.batch);
- decimal64_batch != nullptr) {
- view.decimal_scale = decimal64_batch->scale;
- for (size_t row = 0; row < output_rows; ++row) {
- Int128 value = 0;
- const auto source_row = orc_source_row_at(row,
orc_view.selected_rows);
- if (!orc_row_is_null(*orc_view.batch, source_row)) {
-
RETURN_IF_ERROR(scale_decimal_value(decimal64_batch->values[source_row],
- decimal64_batch->scale,
target_scale, &value));
- }
- fill_decimal_big_endian_value(value, &decimal_values[row]);
- binary_values.emplace_back(reinterpret_cast<const
char*>(decimal_values[row].data()),
- decimal_values[row].size());
- }
- view.binary_values = &binary_values;
- RETURN_IF_ERROR(read_decoded_values(serde, column, &view));
- return Status::OK();
- }
-
- const auto* decimal128_batch =
- dynamic_cast<const ::orc::Decimal128VectorBatch*>(orc_view.batch);
- if (decimal128_batch == nullptr) {
- return Status::InternalError("Unexpected ORC decimal batch type {}",
- orc_view.batch->toString());
- }
- view.decimal_scale = decimal128_batch->scale;
- for (size_t row = 0; row < output_rows; ++row) {
- Int128 value = 0;
- const auto source_row = orc_source_row_at(row, orc_view.selected_rows);
- if (!orc_row_is_null(*orc_view.batch, source_row)) {
-
RETURN_IF_ERROR(scale_decimal_value(to_int128(decimal128_batch->values[source_row]),
- decimal128_batch->scale,
target_scale, &value));
- }
- fill_decimal_big_endian_value(value, &decimal_values[row]);
- binary_values.emplace_back(reinterpret_cast<const
char*>(decimal_values[row].data()),
- decimal_values[row].size());
- }
- view.binary_values = &binary_values;
- RETURN_IF_ERROR(read_decoded_values(serde, column, &view));
- return Status::OK();
-}
-
-Status append_orc_offsets(ColumnArray::Offsets64& doris_offsets,
- const ::orc::DataBuffer<int64_t>& orc_offsets,
size_t rows,
- size_t* element_size, const std::vector<size_t>*
selected_rows,
- std::vector<size_t>* element_selection) {
- DORIS_CHECK(element_size != nullptr);
- if (selected_rows != nullptr) {
- DORIS_CHECK(element_selection != nullptr);
- const auto prev_offset = doris_offsets.empty() ? 0 :
doris_offsets.back();
- ColumnArray::Offset64 current_offset = prev_offset;
- element_selection->clear();
- for (size_t row = 0; row < selected_rows->size(); ++row) {
- const auto source_row = (*selected_rows)[row];
- DORIS_CHECK(source_row < rows);
- const auto begin_offset = orc_offsets[source_row];
- const auto end_offset = orc_offsets[source_row + 1];
- if (end_offset < begin_offset) {
- return Status::Corruption("Invalid ORC offsets");
- }
- const auto delta = static_cast<size_t>(end_offset - begin_offset);
- for (size_t element_idx = 0; element_idx < delta; ++element_idx) {
- element_selection->push_back(static_cast<size_t>(begin_offset)
+ element_idx);
- }
- current_offset += static_cast<ColumnArray::Offset64>(delta);
- doris_offsets.push_back(current_offset);
- }
- *element_size = element_selection->size();
- return Status::OK();
- }
-
- const auto prev_offset = doris_offsets.empty() ? 0 : doris_offsets.back();
- const auto base_offset = orc_offsets[0];
- for (size_t idx = 1; idx <= rows; ++idx) {
- const auto delta = orc_offsets[idx] - base_offset;
- if (delta < 0) {
- return Status::Corruption("Invalid ORC offsets");
- }
- doris_offsets.push_back(prev_offset +
static_cast<ColumnArray::Offset64>(delta));
- }
- const auto total_delta = orc_offsets[rows] - base_offset;
- if (total_delta < 0) {
- return Status::Corruption("Invalid ORC offsets");
- }
- *element_size = static_cast<size_t>(total_delta);
- return Status::OK();
-}
-
-int64_t find_struct_child_index(const ::orc::Type& type, const std::string&
field_name) {
- DORIS_CHECK(type.getKind() == ::orc::TypeKind::STRUCT);
- for (uint64_t child_idx = 0; child_idx < type.getSubtypeCount();
++child_idx) {
- if (type.getFieldName(child_idx) == field_name) {
- return static_cast<int64_t>(child_idx);
- }
- }
- return -1;
-}
-
-struct RoundedOrcTimestamp {
- int64_t seconds;
- uint64_t microseconds;
- bool carry;
-};
-
-Status round_orc_timestamp_to_microseconds(int64_t seconds, int64_t
nanoseconds,
- RoundedOrcTimestamp* result) {
- constexpr int64_t NANOS_PER_SECOND = 1000000000;
- constexpr int64_t NANOS_PER_MICROSECOND = 1000;
- constexpr int64_t MICROS_PER_SECOND = 1000000;
- DORIS_CHECK(result != nullptr);
- DORIS_CHECK(nanoseconds >= 0 && nanoseconds < NANOS_PER_SECOND);
- // Doris stores six fractional digits, so use half-up rounding and carry
999999500ns into the
- // next second instead of silently truncating the ORC value.
- const auto rounded_microseconds =
- (nanoseconds + NANOS_PER_MICROSECOND / 2) / NANOS_PER_MICROSECOND;
- // Validate the carry here, but validate Doris' calendar range after
timezone conversion:
- // a valid year-zero local timestamp may have a UTC epoch before year zero.
- if (__builtin_add_overflow(seconds, rounded_microseconds /
MICROS_PER_SECOND,
- &result->seconds)) {
- return Status::DataQualityError("ORC timestamp overflows after
microsecond rounding");
- }
- result->microseconds = cast_set<uint64_t>(rounded_microseconds %
MICROS_PER_SECOND);
- result->carry = rounded_microseconds >= MICROS_PER_SECOND;
- return Status::OK();
-}
-
-Status decode_timestamp_orc_values(IColumn& nested_column, const
OrcDecodedColumnView& orc_view,
- const cctz::time_zone& timezone) {
- const auto* orc_batch = dynamic_cast<const
::orc::TimestampVectorBatch*>(orc_view.batch);
- if (orc_batch == nullptr) {
- return Status::InternalError("Unexpected ORC timestamp batch type {}",
- orc_view.batch->toString());
- }
- auto& data = assert_cast<ColumnDateTimeV2&>(nested_column).get_data();
- const size_t old_data_size = data.size();
- const auto output_rows = orc_decode_row_count(orc_view.rows,
orc_view.selected_rows);
- data.resize(old_data_size + output_rows);
- for (size_t row = 0; row < output_rows; ++row) {
- const auto source_row = orc_source_row_at(row, orc_view.selected_rows);
- if (orc_row_is_null(*orc_view.batch, source_row)) {
- data[old_data_size + row] = DateV2Value<DateTimeV2ValueType> {};
- continue;
- }
- auto& value =
-
reinterpret_cast<DateV2Value<DateTimeV2ValueType>&>(data[old_data_size + row]);
- RoundedOrcTimestamp timestamp;
- auto status = round_orc_timestamp_to_microseconds(
- orc_batch->data[source_row],
orc_batch->nanoseconds[source_row], ×tamp);
- if (!status.ok()) {
- data.resize(old_data_size);
- return status;
- }
- value.from_unixtime(orc_batch->data[source_row], timezone);
- if (!value.is_valid_date()) {
- data.resize(old_data_size);
- return Status::DataQualityError(
- "Decoded ORC timestamp is outside the target timezone
range");
- }
- value.set_microsecond(timestamp.microseconds);
- // Plain ORC TIMESTAMP is a civil value. Carry after timezone
conversion so a fractional
- // round does not jump backward or skip an hour at a daylight-saving
transition.
- if (timestamp.carry &&
- !value.date_add_interval<TimeUnit::SECOND>(TimeInterval
{TimeUnit::SECOND, 1, false})) {
- data.resize(old_data_size);
- return Status::DataQualityError(
- "Decoded ORC timestamp is outside the target timezone
range");
- }
- }
- return Status::OK();
-}
-
-Status decode_timestamp_tz_orc_values(IColumn& nested_column,
- const OrcDecodedColumnView& orc_view) {
- const auto* orc_batch = dynamic_cast<const
::orc::TimestampVectorBatch*>(orc_view.batch);
- if (orc_batch == nullptr) {
- return Status::InternalError("Unexpected ORC timestamp batch type {}",
- orc_view.batch->toString());
- }
- auto& data = assert_cast<ColumnTimeStampTz&>(nested_column).get_data();
- const size_t old_data_size = data.size();
- const auto output_rows = orc_decode_row_count(orc_view.rows,
orc_view.selected_rows);
- data.resize(old_data_size + output_rows);
- static const auto utc_time_zone = cctz::utc_time_zone();
- for (size_t row = 0; row < output_rows; ++row) {
- const auto source_row = orc_source_row_at(row, orc_view.selected_rows);
- if (orc_row_is_null(*orc_view.batch, source_row)) {
- data[old_data_size + row] = TimestampTzValue {};
- continue;
- }
- auto& value = data[old_data_size + row];
- RoundedOrcTimestamp timestamp;
- auto status = round_orc_timestamp_to_microseconds(
- orc_batch->data[source_row],
orc_batch->nanoseconds[source_row], ×tamp);
- if (!status.ok()) {
- data.resize(old_data_size);
- return status;
- }
- value.from_unixtime(timestamp.seconds, utc_time_zone);
- value.set_microsecond(timestamp.microseconds);
- if (!value.is_valid_date()) {
- data.resize(old_data_size);
- return Status::DataQualityError(
- "Decoded ORC TIMESTAMPTZ is outside the Doris 0000-9999
range");
- }
- }
- return Status::OK();
-}
-
-OrcDecodedColumnView make_child_orc_view(const OrcDecodedColumnView&
parent_view,
- const ::orc::Type* file_type,
- const ::orc::Type* selected_type,
- const ::orc::ColumnVectorBatch*
batch, size_t rows,
- const std::vector<size_t>*
selected_rows) {
- OrcDecodedColumnView child_view = parent_view;
- child_view.file_type = file_type;
- child_view.selected_type = selected_type;
- child_view.batch = batch;
- child_view.rows = rows;
- child_view.selected_rows = selected_rows;
- return child_view;
-}
-
-Status read_orc_child_column(const DataTypeSerDeSPtr& child_serde,
MutableColumnPtr& child_column,
- const OrcDecodedColumnView& child_view) {
- DORIS_CHECK(child_serde != nullptr);
- RETURN_IF_ERROR(child_serde->read_column_from_orc(*child_column,
child_view));
- return Status::OK();
-}
-
-Status decode_list_orc_values(const DataTypeSerDeSPtr& nested_serde, IColumn&
nested_column,
- const OrcDecodedColumnView& orc_view) {
- const auto* orc_list = dynamic_cast<const
::orc::ListVectorBatch*>(orc_view.batch);
- if (orc_list == nullptr) {
- return Status::InternalError("Unexpected ORC list batch type {}",
- orc_view.batch->toString());
- }
- DORIS_CHECK(orc_view.file_type != nullptr);
- DORIS_CHECK(orc_view.selected_type != nullptr);
- DORIS_CHECK(orc_view.file_type->getSubtypeCount() == 1);
- DORIS_CHECK(orc_view.selected_type->getSubtypeCount() == 1);
- DORIS_CHECK(orc_list->elements != nullptr);
- const auto* file_element_type = orc_view.file_type->getSubtype(0);
- const auto* selected_element_type = orc_view.selected_type->getSubtype(0);
- DORIS_CHECK(file_element_type != nullptr);
- DORIS_CHECK(selected_element_type != nullptr);
-
- auto& array_column = assert_cast<ColumnArray&>(nested_column);
- size_t element_size = 0;
- std::vector<size_t> element_selection;
- RETURN_IF_ERROR(append_orc_offsets(array_column.get_offsets(),
orc_list->offsets, orc_view.rows,
- &element_size, orc_view.selected_rows,
&element_selection));
- auto element_column = array_column.get_data_ptr()->assert_mutable();
- const auto child_rows = orc_view.selected_rows == nullptr
- ? element_size
- :
static_cast<size_t>(orc_list->elements->numElements);
- const auto* child_selection = orc_view.selected_rows == nullptr ? nullptr
: &element_selection;
- auto child_view = make_child_orc_view(orc_view, file_element_type,
selected_element_type,
- orc_list->elements.get(),
child_rows, child_selection);
- RETURN_IF_ERROR(read_orc_child_column(nested_serde, element_column,
child_view));
- array_column.get_data_ptr() = std::move(element_column);
- return Status::OK();
-}
-
-Status decode_map_orc_values(const DataTypeSerDeSPtr& key_serde,
- const DataTypeSerDeSPtr& value_serde, IColumn&
nested_column,
- const OrcDecodedColumnView& orc_view) {
- const auto* orc_map = dynamic_cast<const
::orc::MapVectorBatch*>(orc_view.batch);
- if (orc_map == nullptr) {
- return Status::InternalError("Unexpected ORC map batch type {}",
- orc_view.batch->toString());
- }
- DORIS_CHECK(orc_view.file_type != nullptr);
- DORIS_CHECK(orc_view.selected_type != nullptr);
- DORIS_CHECK(orc_view.file_type->getSubtypeCount() == 2);
- DORIS_CHECK(orc_view.selected_type->getSubtypeCount() == 2);
- DORIS_CHECK(orc_map->keys != nullptr);
- DORIS_CHECK(orc_map->elements != nullptr);
- auto& map_column = assert_cast<ColumnMap&>(nested_column);
- size_t element_size = 0;
- std::vector<size_t> element_selection;
- RETURN_IF_ERROR(append_orc_offsets(map_column.get_offsets(),
orc_map->offsets, orc_view.rows,
- &element_size, orc_view.selected_rows,
&element_selection));
-
- const auto* file_key_type = orc_view.file_type->getSubtype(0);
- const auto* selected_key_type = orc_view.selected_type->getSubtype(0);
- DORIS_CHECK(file_key_type != nullptr);
- DORIS_CHECK(selected_key_type != nullptr);
- const auto child_rows = orc_view.selected_rows == nullptr
- ? element_size
- :
static_cast<size_t>(orc_map->keys->numElements);
- const auto* child_selection = orc_view.selected_rows == nullptr ? nullptr
: &element_selection;
- auto key_column = map_column.get_keys_ptr()->assert_mutable();
- auto key_view = make_child_orc_view(orc_view, file_key_type,
selected_key_type,
- orc_map->keys.get(), child_rows,
child_selection);
- RETURN_IF_ERROR(read_orc_child_column(key_serde, key_column, key_view));
- map_column.get_keys_ptr() = std::move(key_column);
-
- const auto* file_value_type = orc_view.file_type->getSubtype(1);
- const auto* selected_value_type = orc_view.selected_type->getSubtype(1);
- DORIS_CHECK(file_value_type != nullptr);
- DORIS_CHECK(selected_value_type != nullptr);
- auto value_column = map_column.get_values_ptr()->assert_mutable();
- auto value_view = make_child_orc_view(
- orc_view, file_value_type, selected_value_type,
orc_map->elements.get(),
- orc_view.selected_rows == nullptr ? element_size
- :
static_cast<size_t>(orc_map->elements->numElements),
- child_selection);
- RETURN_IF_ERROR(read_orc_child_column(value_serde, value_column,
value_view));
- map_column.get_values_ptr() = std::move(value_column);
- return Status::OK();
-}
-
-Status decode_struct_orc_values(const DataTypeSerDeSPtrs& elem_serdes_ptrs,
IColumn& nested_column,
- const OrcDecodedColumnView& orc_view) {
- const auto* orc_struct = dynamic_cast<const
::orc::StructVectorBatch*>(orc_view.batch);
- if (orc_struct == nullptr) {
- return Status::InternalError("Unexpected ORC struct batch type {}",
- orc_view.batch->toString());
- }
- DORIS_CHECK(orc_view.file_type != nullptr);
- DORIS_CHECK(orc_view.selected_type != nullptr);
- DORIS_CHECK(orc_view.selected_type->getSubtypeCount() ==
orc_struct->fields.size());
- auto& struct_column = assert_cast<ColumnStruct&>(nested_column);
- DORIS_CHECK(struct_column.tuple_size() ==
orc_view.selected_type->getSubtypeCount());
- DORIS_CHECK(elem_serdes_ptrs.size() ==
orc_view.selected_type->getSubtypeCount());
-
- for (uint64_t selected_idx = 0; selected_idx <
orc_view.selected_type->getSubtypeCount();
- ++selected_idx) {
- const auto field_name =
orc_view.selected_type->getFieldName(selected_idx);
- const auto file_child_idx =
find_struct_child_index(*orc_view.file_type, field_name);
- if (file_child_idx < 0) {
- return Status::InternalError("Selected ORC field {} is not in file
struct", field_name);
- }
- const auto* file_child_type =
-
orc_view.file_type->getSubtype(static_cast<uint64_t>(file_child_idx));
- const auto* selected_child_type =
orc_view.selected_type->getSubtype(selected_idx);
- DORIS_CHECK(file_child_type != nullptr);
- DORIS_CHECK(selected_child_type != nullptr);
- DORIS_CHECK(selected_idx < orc_struct->fields.size());
- auto child_column =
-
struct_column.get_column_ptr(static_cast<size_t>(selected_idx))->assert_mutable();
- auto child_view = make_child_orc_view(orc_view, file_child_type,
selected_child_type,
-
orc_struct->fields[selected_idx], orc_view.rows,
- orc_view.selected_rows);
- RETURN_IF_ERROR(
- read_orc_child_column(elem_serdes_ptrs[selected_idx],
child_column, child_view));
- struct_column.get_column_ptr(static_cast<size_t>(selected_idx)) =
std::move(child_column);
- }
- return Status::OK();
-}
-
-} // namespace
-
DataTypeSerDe::~DataTypeSerDe() = default;
bool decoded_column_view_can_null_on_conversion_failure(const
DecodedColumnView& view) {
@@ -799,213 +105,6 @@ Status DataTypeSerDe::read_column_from_orc(IColumn&
column,
return Status::NotSupported("read_column_from_orc is not supported for
{}", get_name());
}
-Status DataTypeNullableSerDe::read_column_from_orc(IColumn& column,
- const OrcDecodedColumnView&
view) const {
- DORIS_CHECK(view.file_type != nullptr);
- DORIS_CHECK(view.selected_type != nullptr);
- DORIS_CHECK(view.batch != nullptr);
- DORIS_CHECK(view.file_type->getKind() == view.selected_type->getKind());
- auto& nullable_column = assert_cast<ColumnNullable&>(column);
- const auto output_rows = orc_decode_row_count(view.rows,
view.selected_rows);
- if (output_rows == 0) {
- return Status::OK();
- }
-
- auto& null_map = nullable_column.get_null_map_data();
- const auto old_null_map_size = null_map.size();
- auto& nested_column = nullable_column.get_nested_column();
- const auto old_nested_size = nested_column.size();
- append_orc_null_map(*view.batch, view.rows, view.selected_rows, &null_map);
- auto st = nested_serde->read_column_from_orc(nested_column, view);
- if (!st.ok()) {
- null_map.resize(old_null_map_size);
- nested_column.resize(old_nested_size);
- }
- return st;
-}
-
-template <PrimitiveType T>
-Status DataTypeNumberSerDe<T>::read_column_from_orc(IColumn& column,
- const
OrcDecodedColumnView& view) const {
- DORIS_CHECK(view.file_type != nullptr);
- DORIS_CHECK(view.batch != nullptr);
- if (orc_decode_row_count(view.rows, view.selected_rows) == 0) {
- return Status::OK();
- }
-
- if constexpr (T == TYPE_BOOLEAN) {
- DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::BOOLEAN);
- return decode_boolean_orc_values(*this, column, view);
- } else if constexpr (T == TYPE_TINYINT || T == TYPE_SMALLINT || T ==
TYPE_INT ||
- T == TYPE_BIGINT) {
- if constexpr (T == TYPE_TINYINT) {
- DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::BYTE);
- } else if constexpr (T == TYPE_SMALLINT) {
- DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::SHORT);
- } else if constexpr (T == TYPE_INT) {
- DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::INT);
- } else {
- DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::LONG);
- }
- return decode_fixed_orc_values<::orc::LongVectorBatch, int64_t>(*this,
column, view,
-
DecodedValueKind::INT64);
- } else if constexpr (T == TYPE_FLOAT) {
- DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::FLOAT);
- return decode_float_orc_values(*this, column, view);
- } else if constexpr (T == TYPE_DOUBLE) {
- DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::DOUBLE);
- return decode_fixed_orc_values<::orc::DoubleVectorBatch,
double>(*this, column, view,
-
DecodedValueKind::DOUBLE);
- }
- return DataTypeSerDe::read_column_from_orc(column, view);
-}
-
-template <typename ColumnType>
-Status DataTypeStringSerDeBase<ColumnType>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const {
- DORIS_CHECK(view.file_type != nullptr);
- DORIS_CHECK(view.batch != nullptr);
- const auto kind = view.file_type->getKind();
- DORIS_CHECK(kind == ::orc::TypeKind::STRING || kind ==
::orc::TypeKind::BINARY ||
- kind == ::orc::TypeKind::VARCHAR || kind ==
::orc::TypeKind::CHAR);
- if (orc_decode_row_count(view.rows, view.selected_rows) == 0) {
- return Status::OK();
- }
- return decode_string_orc_values(*this, column, view);
-}
-
-template <PrimitiveType T>
-Status DataTypeDecimalSerDe<T>::read_column_from_orc(IColumn& column,
- const
OrcDecodedColumnView& view) const {
- DORIS_CHECK(view.file_type != nullptr);
- DORIS_CHECK(view.batch != nullptr);
- DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::DECIMAL);
- if (orc_decode_row_count(view.rows, view.selected_rows) == 0) {
- return Status::OK();
- }
- return decode_decimal_orc_values(*this, column, view,
cast_set<int32_t>(scale));
-}
-
-Status DataTypeDateV2SerDe::read_column_from_orc(IColumn& column,
- const OrcDecodedColumnView&
view) const {
- DORIS_CHECK(view.file_type != nullptr);
- DORIS_CHECK(view.batch != nullptr);
- DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::DATE);
- if (orc_decode_row_count(view.rows, view.selected_rows) == 0) {
- return Status::OK();
- }
- return decode_date_orc_values(*this, column, view);
-}
-
-Status DataTypeDateTimeV2SerDe::read_column_from_orc(IColumn& column,
- const
OrcDecodedColumnView& view) const {
- DORIS_CHECK(view.file_type != nullptr);
- DORIS_CHECK(view.batch != nullptr);
- const auto kind = view.file_type->getKind();
- DORIS_CHECK(kind == ::orc::TypeKind::TIMESTAMP || kind ==
::orc::TypeKind::TIMESTAMP_INSTANT);
- DORIS_CHECK(view.timezone != nullptr);
- if (orc_decode_row_count(view.rows, view.selected_rows) == 0) {
- return Status::OK();
- }
- return decode_timestamp_orc_values(column, view, *view.timezone);
-}
-
-Status DataTypeTimeStampTzSerDe::read_column_from_orc(IColumn& column,
- const
OrcDecodedColumnView& view) const {
- DORIS_CHECK(view.file_type != nullptr);
- DORIS_CHECK(view.batch != nullptr);
- DORIS_CHECK(view.file_type->getKind() ==
::orc::TypeKind::TIMESTAMP_INSTANT);
- DORIS_CHECK(view.enable_mapping_timestamp_tz);
- if (orc_decode_row_count(view.rows, view.selected_rows) == 0) {
- return Status::OK();
- }
- return decode_timestamp_tz_orc_values(column, view);
-}
-
-Status DataTypeArraySerDe::read_column_from_orc(IColumn& column,
- const OrcDecodedColumnView&
view) const {
- DORIS_CHECK(view.file_type != nullptr);
- DORIS_CHECK(view.batch != nullptr);
- DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::LIST);
- if (orc_decode_row_count(view.rows, view.selected_rows) == 0) {
- return Status::OK();
- }
- return decode_list_orc_values(nested_serde, column, view);
-}
-
-Status DataTypeMapSerDe::read_column_from_orc(IColumn& column,
- const OrcDecodedColumnView&
view) const {
- DORIS_CHECK(view.file_type != nullptr);
- DORIS_CHECK(view.batch != nullptr);
- DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::MAP);
- if (orc_decode_row_count(view.rows, view.selected_rows) == 0) {
- return Status::OK();
- }
- return decode_map_orc_values(key_serde, value_serde, column, view);
-}
-
-Status DataTypeStructSerDe::read_column_from_orc(IColumn& column,
- const OrcDecodedColumnView&
view) const {
- DORIS_CHECK(view.file_type != nullptr);
- DORIS_CHECK(view.batch != nullptr);
- DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::STRUCT);
- if (orc_decode_row_count(view.rows, view.selected_rows) == 0) {
- return Status::OK();
- }
- return decode_struct_orc_values(elem_serdes_ptrs, column, view);
-}
-
-template Status DataTypeNumberSerDe<TYPE_BOOLEAN>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-template Status DataTypeNumberSerDe<TYPE_TINYINT>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-template Status DataTypeNumberSerDe<TYPE_SMALLINT>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-template Status DataTypeNumberSerDe<TYPE_INT>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-template Status DataTypeNumberSerDe<TYPE_BIGINT>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-template Status DataTypeNumberSerDe<TYPE_LARGEINT>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-template Status DataTypeNumberSerDe<TYPE_FLOAT>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-template Status DataTypeNumberSerDe<TYPE_DOUBLE>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-template Status DataTypeNumberSerDe<TYPE_DATE>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-template Status DataTypeNumberSerDe<TYPE_DATEV2>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-template Status DataTypeNumberSerDe<TYPE_DATETIME>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-template Status DataTypeNumberSerDe<TYPE_DATETIMEV2>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-template Status DataTypeNumberSerDe<TYPE_IPV4>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-template Status DataTypeNumberSerDe<TYPE_IPV6>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-template Status DataTypeNumberSerDe<TYPE_TIMEV2>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-template Status DataTypeNumberSerDe<TYPE_TIMESTAMPTZ>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-
-template Status DataTypeStringSerDeBase<ColumnString>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-template Status DataTypeStringSerDeBase<ColumnString64>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-template Status
DataTypeStringSerDeBase<ColumnFixedLengthObject>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-
-template Status DataTypeDecimalSerDe<TYPE_DECIMAL32>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-template Status DataTypeDecimalSerDe<TYPE_DECIMAL64>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-template Status DataTypeDecimalSerDe<TYPE_DECIMAL128I>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-template Status DataTypeDecimalSerDe<TYPE_DECIMALV2>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-template Status DataTypeDecimalSerDe<TYPE_DECIMAL256>::read_column_from_orc(
- IColumn& column, const OrcDecodedColumnView& view) const;
-
Status DataTypeSerDe::read_field_from_decoded_value(const IDataType&
data_type, Field* field,
const DecodedColumnView&
view) const {
DORIS_CHECK(field != nullptr);
diff --git a/be/src/core/data_type_serde/data_type_string_serde.cpp
b/be/src/core/data_type_serde/data_type_string_serde.cpp
index 24f7215d42e..0fc469743e4 100644
--- a/be/src/core/data_type_serde/data_type_string_serde.cpp
+++ b/be/src/core/data_type_serde/data_type_string_serde.cpp
@@ -28,6 +28,7 @@
#include "core/data_type/define_primitive_type.h"
#include "core/data_type_serde/arrow_validation.h"
#include "core/data_type_serde/decoded_column_view.h"
+#include "core/data_type_serde/orc_serde_utils.h"
#include "core/data_type_serde/parquet_decode_source.h"
#include "util/jsonb_document_cast.h"
#include "util/jsonb_utils.h"
@@ -36,6 +37,90 @@
namespace doris {
namespace {
+size_t trim_right_spaces(const char* value, size_t length) {
+ while (length > 0 && value[length - 1] == ' ') {
+ --length;
+ }
+ return length;
+}
+
+Status append_orc_string_ref(const ::orc::Type& file_type, const char* data,
int64_t length,
+ std::vector<StringRef>& binary_values) {
+ if (length < 0) {
+ return Status::Corruption("Invalid negative ORC string length {}",
length);
+ }
+ auto value_length = static_cast<size_t>(length);
+ if (file_type.getKind() == ::orc::TypeKind::CHAR) {
+ value_length = trim_right_spaces(data, value_length);
+ }
+ binary_values.emplace_back(value_length == 0 ? "" : data, value_length);
+ return Status::OK();
+}
+
+Status decode_string_orc_values(const DataTypeSerDe& serde, IColumn& column,
+ const OrcDecodedColumnView& orc_view) {
+ DORIS_CHECK(orc_view.file_type != nullptr);
+ if (const auto* encoded_batch =
+ dynamic_cast<const
::orc::EncodedStringVectorBatch*>(orc_view.batch);
+ encoded_batch != nullptr && encoded_batch->isEncoded) {
+ if (encoded_batch->dictionary == nullptr) {
+ return Status::InternalError("Encoded ORC string batch has no
dictionary");
+ }
+ auto view = orc_serde_utils::make_orc_decoded_view(orc_view,
DecodedValueKind::BINARY);
+ NullMap null_map;
+ orc_serde_utils::fill_orc_decoded_null_map(*orc_view.batch,
orc_view.rows,
+ orc_view.selected_rows,
&null_map);
+ view.null_map = null_map.empty() ? nullptr : null_map.data();
+ const auto output_rows =
+ orc_serde_utils::orc_decode_row_count(orc_view.rows,
orc_view.selected_rows);
+ std::vector<StringRef> binary_values;
+ binary_values.reserve(output_rows);
+ for (size_t row = 0; row < output_rows; ++row) {
+ const auto source_row = orc_serde_utils::orc_source_row_at(row,
orc_view.selected_rows);
+ if (orc_serde_utils::orc_row_is_null(*orc_view.batch, source_row))
{
+ binary_values.emplace_back("", 0);
+ continue;
+ }
+ char* data = nullptr;
+ int64_t length = 0;
+
encoded_batch->dictionary->getValueByIndex(encoded_batch->index[source_row],
data,
+ length);
+ RETURN_IF_ERROR(
+ append_orc_string_ref(*orc_view.file_type, data, length,
binary_values));
+ }
+ view.binary_values = &binary_values;
+ RETURN_IF_ERROR(orc_serde_utils::read_decoded_values(serde, column,
&view));
+ return Status::OK();
+ }
+
+ const auto* orc_batch = dynamic_cast<const
::orc::StringVectorBatch*>(orc_view.batch);
+ if (orc_batch == nullptr) {
+ return Status::InternalError("Unexpected ORC string batch type {}",
+ orc_view.batch->toString());
+ }
+ auto view = orc_serde_utils::make_orc_decoded_view(orc_view,
DecodedValueKind::BINARY);
+ NullMap null_map;
+ orc_serde_utils::fill_orc_decoded_null_map(*orc_view.batch, orc_view.rows,
+ orc_view.selected_rows,
&null_map);
+ view.null_map = null_map.empty() ? nullptr : null_map.data();
+ const auto output_rows =
+ orc_serde_utils::orc_decode_row_count(orc_view.rows,
orc_view.selected_rows);
+ std::vector<StringRef> binary_values;
+ binary_values.reserve(output_rows);
+ for (size_t row = 0; row < output_rows; ++row) {
+ const auto source_row = orc_serde_utils::orc_source_row_at(row,
orc_view.selected_rows);
+ if (orc_serde_utils::orc_row_is_null(*orc_view.batch, source_row)) {
+ binary_values.emplace_back("", 0);
+ continue;
+ }
+ RETURN_IF_ERROR(append_orc_string_ref(*orc_view.file_type,
orc_batch->data[source_row],
+ orc_batch->length[source_row],
binary_values));
+ }
+ view.binary_values = &binary_values;
+ RETURN_IF_ERROR(orc_serde_utils::read_decoded_values(serde, column,
&view));
+ return Status::OK();
+}
+
template <typename ColumnType>
Status read_string_decoded_values(IColumn& column, const DecodedColumnView&
view) {
if (view.binary_values == nullptr &&
decoded_column_view_has_non_null_value(view)) {
@@ -772,6 +857,20 @@ Status
DataTypeStringSerDeBase<ColumnType>::from_olap_string(const std::string&
return Status::OK();
}
+template <typename ColumnType>
+Status DataTypeStringSerDeBase<ColumnType>::read_column_from_orc(
+ IColumn& column, const OrcDecodedColumnView& view) const {
+ DORIS_CHECK(view.file_type != nullptr);
+ DORIS_CHECK(view.batch != nullptr);
+ const auto kind = view.file_type->getKind();
+ DORIS_CHECK(kind == ::orc::TypeKind::STRING || kind ==
::orc::TypeKind::BINARY ||
+ kind == ::orc::TypeKind::VARCHAR || kind ==
::orc::TypeKind::CHAR);
+ if (orc_serde_utils::orc_decode_row_count(view.rows, view.selected_rows)
== 0) {
+ return Status::OK();
+ }
+ return decode_string_orc_values(*this, column, view);
+}
+
template class DataTypeStringSerDeBase<ColumnString>;
template class DataTypeStringSerDeBase<ColumnString64>;
template class DataTypeStringSerDeBase<ColumnFixedLengthObject>;
diff --git a/be/src/core/data_type_serde/data_type_struct_serde.cpp
b/be/src/core/data_type_serde/data_type_struct_serde.cpp
index 15f40a54aa5..8d7cf1fa30e 100644
--- a/be/src/core/data_type_serde/data_type_struct_serde.cpp
+++ b/be/src/core/data_type_serde/data_type_struct_serde.cpp
@@ -28,6 +28,7 @@
#include "core/data_type_serde/arrow_validation.h"
#include "core/data_type_serde/complex_type_deserialize_util.h"
#include "core/data_type_serde/data_type_serde.h"
+#include "core/data_type_serde/orc_serde_utils.h"
#include "core/string_ref.h"
#include "util/jsonb_document.h"
#include "util/jsonb_writer.h"
@@ -36,6 +37,58 @@ namespace doris {
class Arena;
+namespace {
+
+int64_t find_struct_child_index(const ::orc::Type& type, const std::string&
field_name) {
+ DORIS_CHECK(type.getKind() == ::orc::TypeKind::STRUCT);
+ for (uint64_t child_idx = 0; child_idx < type.getSubtypeCount();
++child_idx) {
+ if (type.getFieldName(child_idx) == field_name) {
+ return static_cast<int64_t>(child_idx);
+ }
+ }
+ return -1;
+}
+
+Status decode_struct_orc_values(const DataTypeSerDeSPtrs& elem_serdes_ptrs,
IColumn& nested_column,
+ const OrcDecodedColumnView& orc_view) {
+ const auto* orc_struct = dynamic_cast<const
::orc::StructVectorBatch*>(orc_view.batch);
+ if (orc_struct == nullptr) {
+ return Status::InternalError("Unexpected ORC struct batch type {}",
+ orc_view.batch->toString());
+ }
+ DORIS_CHECK(orc_view.file_type != nullptr);
+ DORIS_CHECK(orc_view.selected_type != nullptr);
+ DORIS_CHECK(orc_view.selected_type->getSubtypeCount() ==
orc_struct->fields.size());
+ auto& struct_column = assert_cast<ColumnStruct&>(nested_column);
+ DORIS_CHECK(struct_column.tuple_size() ==
orc_view.selected_type->getSubtypeCount());
+ DORIS_CHECK(elem_serdes_ptrs.size() ==
orc_view.selected_type->getSubtypeCount());
+ for (uint64_t selected_idx = 0; selected_idx <
orc_view.selected_type->getSubtypeCount();
+ ++selected_idx) {
+ const auto field_name =
orc_view.selected_type->getFieldName(selected_idx);
+ const auto file_child_idx =
find_struct_child_index(*orc_view.file_type, field_name);
+ if (file_child_idx < 0) {
+ return Status::InternalError("Selected ORC field {} is not in file
struct", field_name);
+ }
+ const auto* file_child_type =
+
orc_view.file_type->getSubtype(static_cast<uint64_t>(file_child_idx));
+ const auto* selected_child_type =
orc_view.selected_type->getSubtype(selected_idx);
+ DORIS_CHECK(file_child_type != nullptr);
+ DORIS_CHECK(selected_child_type != nullptr);
+ DORIS_CHECK(selected_idx < orc_struct->fields.size());
+ auto child_column =
+
struct_column.get_column_ptr(static_cast<size_t>(selected_idx))->assert_mutable();
+ auto child_view = orc_serde_utils::make_child_orc_view(
+ orc_view, file_child_type, selected_child_type,
orc_struct->fields[selected_idx],
+ orc_view.rows, orc_view.selected_rows);
+
RETURN_IF_ERROR(orc_serde_utils::read_orc_child_column(elem_serdes_ptrs[selected_idx],
+ child_column,
child_view));
+ struct_column.get_column_ptr(static_cast<size_t>(selected_idx)) =
std::move(child_column);
+ }
+ return Status::OK();
+}
+
+} // namespace
+
std::string DataTypeStructSerDe::get_name() const {
size_t size = elem_names.size();
std::stringstream s;
@@ -667,4 +720,15 @@ bool DataTypeStructSerDe::write_column_to_hive_text(const
IColumn& column, Buffe
return true;
}
+Status DataTypeStructSerDe::read_column_from_orc(IColumn& column,
+ const OrcDecodedColumnView&
view) const {
+ DORIS_CHECK(view.file_type != nullptr);
+ DORIS_CHECK(view.batch != nullptr);
+ DORIS_CHECK(view.file_type->getKind() == ::orc::TypeKind::STRUCT);
+ if (orc_serde_utils::orc_decode_row_count(view.rows, view.selected_rows)
== 0) {
+ return Status::OK();
+ }
+ return decode_struct_orc_values(elem_serdes_ptrs, column, view);
+}
+
} // namespace doris
diff --git a/be/src/core/data_type_serde/data_type_timestamptz_serde.cpp
b/be/src/core/data_type_serde/data_type_timestamptz_serde.cpp
index e9968e0bc87..1e84fef5b05 100644
--- a/be/src/core/data_type_serde/data_type_timestamptz_serde.cpp
+++ b/be/src/core/data_type_serde/data_type_timestamptz_serde.cpp
@@ -24,6 +24,7 @@
#include "core/data_type/primitive_type.h"
#include "core/data_type_serde/arrow_validation.h"
#include "core/data_type_serde/decoded_column_view.h"
+#include "core/data_type_serde/orc_serde_utils.h"
#include "core/data_type_serde/parquet_decode_source.h"
#include "core/data_type_serde/parquet_timestamp.h"
#include "core/value/timestamptz_value.h"
@@ -35,6 +36,44 @@ namespace doris {
namespace {
+Status decode_timestamp_tz_orc_values(IColumn& nested_column,
+ const OrcDecodedColumnView& orc_view) {
+ const auto* orc_batch = dynamic_cast<const
::orc::TimestampVectorBatch*>(orc_view.batch);
+ if (orc_batch == nullptr) {
+ return Status::InternalError("Unexpected ORC timestamp batch type {}",
+ orc_view.batch->toString());
+ }
+ auto& data = assert_cast<ColumnTimeStampTz&>(nested_column).get_data();
+ const size_t old_data_size = data.size();
+ const auto output_rows =
+ orc_serde_utils::orc_decode_row_count(orc_view.rows,
orc_view.selected_rows);
+ data.resize(old_data_size + output_rows);
+ static const auto utc_time_zone = cctz::utc_time_zone();
+ for (size_t row = 0; row < output_rows; ++row) {
+ const auto source_row = orc_serde_utils::orc_source_row_at(row,
orc_view.selected_rows);
+ if (orc_serde_utils::orc_row_is_null(*orc_view.batch, source_row)) {
+ data[old_data_size + row] = TimestampTzValue {};
+ continue;
+ }
+ auto& value = data[old_data_size + row];
+ orc_serde_utils::RoundedOrcTimestamp timestamp;
+ auto status = orc_serde_utils::round_orc_timestamp_to_microseconds(
+ orc_batch->data[source_row],
orc_batch->nanoseconds[source_row], ×tamp);
+ if (!status.ok()) {
+ data.resize(old_data_size);
+ return status;
+ }
+ value.from_unixtime(timestamp.seconds, utc_time_zone);
+ value.set_microsecond(timestamp.microseconds);
+ if (!value.is_valid_date()) {
+ data.resize(old_data_size);
+ return Status::DataQualityError(
+ "Decoded ORC TIMESTAMPTZ is outside the Doris 0000-9999
range");
+ }
+ }
+ return Status::OK();
+}
+
Status append_timestamptz_from_utc_epoch_micros(ColumnTimeStampTz::Container&
data,
int64_t timestamp_micros) {
static constexpr int64_t MICROS_PER_SECOND = 1000000;
@@ -604,4 +643,16 @@ void
DataTypeTimeStampTzSerDe::write_one_cell_to_binary(const IColumn& src_colum
data_ref.size);
}
+Status DataTypeTimeStampTzSerDe::read_column_from_orc(IColumn& column,
+ const
OrcDecodedColumnView& view) const {
+ DORIS_CHECK(view.file_type != nullptr);
+ DORIS_CHECK(view.batch != nullptr);
+ DORIS_CHECK(view.file_type->getKind() ==
::orc::TypeKind::TIMESTAMP_INSTANT);
+ DORIS_CHECK(view.enable_mapping_timestamp_tz);
+ if (orc_serde_utils::orc_decode_row_count(view.rows, view.selected_rows)
== 0) {
+ return Status::OK();
+ }
+ return decode_timestamp_tz_orc_values(column, view);
+}
+
} // namespace doris
diff --git a/be/src/core/data_type_serde/orc_serde_utils.cpp
b/be/src/core/data_type_serde/orc_serde_utils.cpp
new file mode 100644
index 00000000000..5b7c88dbd63
--- /dev/null
+++ b/be/src/core/data_type_serde/orc_serde_utils.cpp
@@ -0,0 +1,161 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+#include "core/data_type_serde/orc_serde_utils.h"
+
+#include "common/cast_set.h"
+#include "common/check.h"
+#include "core/column/column_array.h"
+
+namespace doris::orc_serde_utils {
+
+size_t orc_decode_row_count(size_t rows, const std::vector<size_t>*
selected_rows) {
+ if (selected_rows == nullptr) {
+ return rows;
+ }
+ return selected_rows->size();
+}
+
+size_t orc_source_row_at(size_t row, const std::vector<size_t>* selected_rows)
{
+ if (selected_rows == nullptr) {
+ return row;
+ }
+ return (*selected_rows)[row];
+}
+
+bool orc_row_is_null(const ::orc::ColumnVectorBatch& batch, size_t row) {
+ return batch.hasNulls && !batch.notNull[row];
+}
+
+Status round_orc_timestamp_to_microseconds(int64_t seconds, int64_t
nanoseconds,
+ RoundedOrcTimestamp* result) {
+ constexpr int64_t NANOS_PER_SECOND = 1000000000;
+ constexpr int64_t NANOS_PER_MICROSECOND = 1000;
+ constexpr int64_t MICROS_PER_SECOND = 1000000;
+ DORIS_CHECK(result != nullptr);
+ DORIS_CHECK(nanoseconds >= 0 && nanoseconds < NANOS_PER_SECOND);
+ // Doris stores six fractional digits, so use half-up rounding and carry
999999500ns into the
+ // next second instead of silently truncating the ORC value.
+ const auto rounded_microseconds =
+ (nanoseconds + NANOS_PER_MICROSECOND / 2) / NANOS_PER_MICROSECOND;
+ // Validate the carry here, but validate Doris' calendar range after
timezone conversion:
+ // a valid year-zero local timestamp may have a UTC epoch before year zero.
+ if (__builtin_add_overflow(seconds, rounded_microseconds /
MICROS_PER_SECOND,
+ &result->seconds)) {
+ return Status::DataQualityError("ORC timestamp overflows after
microsecond rounding");
+ }
+ result->microseconds = cast_set<uint64_t>(rounded_microseconds %
MICROS_PER_SECOND);
+ result->carry = rounded_microseconds >= MICROS_PER_SECOND;
+ return Status::OK();
+}
+
+DecodedColumnView make_orc_decoded_view(const OrcDecodedColumnView& orc_view,
+ DecodedValueKind value_kind) {
+ DecodedColumnView view;
+ view.value_kind = value_kind;
+ view.row_count = cast_set<int64_t>(orc_decode_row_count(orc_view.rows,
orc_view.selected_rows));
+ view.timezone = orc_view.timezone;
+ return view;
+}
+
+Status read_decoded_values(const DataTypeSerDe& serde, IColumn& column,
DecodedColumnView* view) {
+ DORIS_CHECK(view != nullptr);
+ RETURN_IF_ERROR(serde.read_column_from_decoded_values(column, *view));
+ return Status::OK();
+}
+
+void fill_orc_decoded_null_map(const ::orc::ColumnVectorBatch& batch, size_t
rows,
+ const std::vector<size_t>* selected_rows,
NullMap* null_map) {
+ DORIS_CHECK(null_map != nullptr);
+ if (!batch.hasNulls) {
+ return;
+ }
+ const auto output_rows = orc_decode_row_count(rows, selected_rows);
+ null_map->resize(output_rows);
+ for (size_t row = 0; row < output_rows; ++row) {
+ (*null_map)[row] = !batch.notNull[orc_source_row_at(row,
selected_rows)];
+ }
+}
+
+Status append_orc_offsets(ColumnArray::Offsets64& doris_offsets,
+ const ::orc::DataBuffer<int64_t>& orc_offsets,
size_t rows,
+ size_t* element_size, const std::vector<size_t>*
selected_rows,
+ std::vector<size_t>* element_selection) {
+ DORIS_CHECK(element_size != nullptr);
+ if (selected_rows != nullptr) {
+ DORIS_CHECK(element_selection != nullptr);
+ const auto prev_offset = doris_offsets.empty() ? 0 :
doris_offsets.back();
+ ColumnArray::Offset64 current_offset = prev_offset;
+ element_selection->clear();
+ for (size_t row = 0; row < selected_rows->size(); ++row) {
+ const auto source_row = (*selected_rows)[row];
+ DORIS_CHECK(source_row < rows);
+ const auto begin_offset = orc_offsets[source_row];
+ const auto end_offset = orc_offsets[source_row + 1];
+ if (end_offset < begin_offset) {
+ return Status::Corruption("Invalid ORC offsets");
+ }
+ const auto delta = static_cast<size_t>(end_offset - begin_offset);
+ for (size_t element_idx = 0; element_idx < delta; ++element_idx) {
+ element_selection->push_back(static_cast<size_t>(begin_offset)
+ element_idx);
+ }
+ current_offset += static_cast<ColumnArray::Offset64>(delta);
+ doris_offsets.push_back(current_offset);
+ }
+ *element_size = element_selection->size();
+ return Status::OK();
+ }
+
+ const auto prev_offset = doris_offsets.empty() ? 0 : doris_offsets.back();
+ const auto base_offset = orc_offsets[0];
+ for (size_t idx = 1; idx <= rows; ++idx) {
+ const auto delta = orc_offsets[idx] - base_offset;
+ if (delta < 0) {
+ return Status::Corruption("Invalid ORC offsets");
+ }
+ doris_offsets.push_back(prev_offset +
static_cast<ColumnArray::Offset64>(delta));
+ }
+ const auto total_delta = orc_offsets[rows] - base_offset;
+ if (total_delta < 0) {
+ return Status::Corruption("Invalid ORC offsets");
+ }
+ *element_size = static_cast<size_t>(total_delta);
+ return Status::OK();
+}
+
+OrcDecodedColumnView make_child_orc_view(const OrcDecodedColumnView&
parent_view,
+ const ::orc::Type* file_type,
+ const ::orc::Type* selected_type,
+ const ::orc::ColumnVectorBatch*
batch, size_t rows,
+ const std::vector<size_t>*
selected_rows) {
+ OrcDecodedColumnView child_view = parent_view;
+ child_view.file_type = file_type;
+ child_view.selected_type = selected_type;
+ child_view.batch = batch;
+ child_view.rows = rows;
+ child_view.selected_rows = selected_rows;
+ return child_view;
+}
+
+Status read_orc_child_column(const DataTypeSerDeSPtr& child_serde,
MutableColumnPtr& child_column,
+ const OrcDecodedColumnView& child_view) {
+ DORIS_CHECK(child_serde != nullptr);
+ RETURN_IF_ERROR(child_serde->read_column_from_orc(*child_column,
child_view));
+ return Status::OK();
+}
+
+} // namespace doris::orc_serde_utils
diff --git a/be/src/core/data_type_serde/orc_serde_utils.h
b/be/src/core/data_type_serde/orc_serde_utils.h
index 0803daab1b0..06b9c4406ee 100644
--- a/be/src/core/data_type_serde/orc_serde_utils.h
+++ b/be/src/core/data_type_serde/orc_serde_utils.h
@@ -19,10 +19,49 @@
#include <cstring>
#include <orc/Vector.hh>
+#include <vector>
#include "core/arena.h"
+#include "core/column/column_array.h"
+#include "core/data_type_serde/data_type_serde.h"
namespace doris {
+namespace orc_serde_utils {
+
+size_t orc_decode_row_count(size_t rows, const std::vector<size_t>*
selected_rows);
+size_t orc_source_row_at(size_t row, const std::vector<size_t>* selected_rows);
+bool orc_row_is_null(const ::orc::ColumnVectorBatch& batch, size_t row);
+
+struct RoundedOrcTimestamp {
+ int64_t seconds;
+ uint64_t microseconds;
+ bool carry;
+};
+
+Status round_orc_timestamp_to_microseconds(int64_t seconds, int64_t
nanoseconds,
+ RoundedOrcTimestamp* result);
+
+DecodedColumnView make_orc_decoded_view(const OrcDecodedColumnView& orc_view,
+ DecodedValueKind value_kind);
+
+Status read_decoded_values(const DataTypeSerDe& serde, IColumn& column,
DecodedColumnView* view);
+
+void fill_orc_decoded_null_map(const ::orc::ColumnVectorBatch& batch, size_t
rows,
+ const std::vector<size_t>* selected_rows,
NullMap* null_map);
+
+Status append_orc_offsets(ColumnArray::Offsets64& doris_offsets,
+ const ::orc::DataBuffer<int64_t>& orc_offsets,
size_t rows,
+ size_t* element_size, const std::vector<size_t>*
selected_rows,
+ std::vector<size_t>* element_selection);
+
+OrcDecodedColumnView make_child_orc_view(const OrcDecodedColumnView&
parent_view,
+ const ::orc::Type* file_type,
+ const ::orc::Type* selected_type,
+ const ::orc::ColumnVectorBatch*
batch, size_t rows,
+ const std::vector<size_t>*
selected_rows);
+
+Status read_orc_child_column(const DataTypeSerDeSPtr& child_serde,
MutableColumnPtr& child_column,
+ const OrcDecodedColumnView& child_view);
inline void copy_orc_string_data_to_arena(orc::ColumnVectorBatch* batch,
Arena& arena) {
if (auto* strings = dynamic_cast<orc::StringVectorBatch*>(batch)) {
@@ -84,4 +123,5 @@ inline void
copy_orc_string_data_to_arena(orc::ColumnVectorBatch* batch, Arena&
}
}
+} // namespace orc_serde_utils
} // namespace doris
diff --git a/be/test/format/transformer/vorc_transformer_test.cpp
b/be/test/format/transformer/vorc_transformer_test.cpp
index ae87c606427..49af5e66244 100644
--- a/be/test/format/transformer/vorc_transformer_test.cpp
+++ b/be/test/format/transformer/vorc_transformer_test.cpp
@@ -435,7 +435,7 @@ TEST(OrcSerdeUtilsTest, CopiesOnlyBorrowedStringData) {
batch.length[1] = borrowed.size();
const size_t used_before_copy = arena.used_size();
- copy_orc_string_data_to_arena(&batch, arena);
+ orc_serde_utils::copy_orc_string_data_to_arena(&batch, arena);
EXPECT_EQ(arena_owned, batch.data[0]);
EXPECT_NE(borrowed.data(), batch.data[1]);
@@ -450,7 +450,7 @@ TEST(OrcSerdeUtilsTest, PreservesEmptyStringAsPresentValue)
{
batch.data[0] = const_cast<char*>("");
batch.length[0] = 0;
- copy_orc_string_data_to_arena(&batch, arena);
+ orc_serde_utils::copy_orc_string_data_to_arena(&batch, arena);
EXPECT_NE(batch.data[0], nullptr);
EXPECT_EQ(batch.length[0], 0);
@@ -471,7 +471,7 @@ TEST(OrcSerdeUtilsTest,
IgnoresBorrowedPayloadForNullString) {
batch.length[1] = 0;
const size_t used_before_copy = arena.used_size();
- copy_orc_string_data_to_arena(&batch, arena);
+ orc_serde_utils::copy_orc_string_data_to_arena(&batch, arena);
EXPECT_EQ(used_before_copy, arena.used_size());
EXPECT_EQ(nullptr, batch.data[0]);
@@ -479,6 +479,27 @@ TEST(OrcSerdeUtilsTest,
IgnoresBorrowedPayloadForNullString) {
EXPECT_NE(nullptr, batch.data[1]);
}
+TEST(OrcSerdeUtilsTest, AppliesSelectedRowsWhenBuildingNullMap) {
+ const std::vector<size_t> selected_rows {3, 1};
+ EXPECT_EQ(2, orc_serde_utils::orc_decode_row_count(4, &selected_rows));
+ EXPECT_EQ(3, orc_serde_utils::orc_source_row_at(0, &selected_rows));
+ EXPECT_EQ(1, orc_serde_utils::orc_source_row_at(1, &selected_rows));
+
+ orc::LongVectorBatch batch(4, *orc::getDefaultPool());
+ batch.numElements = 4;
+ batch.hasNulls = true;
+ batch.notNull[0] = 1;
+ batch.notNull[1] = 0;
+ batch.notNull[2] = 1;
+ batch.notNull[3] = 1;
+
+ NullMap null_map;
+ orc_serde_utils::fill_orc_decoded_null_map(batch, 4, &selected_rows,
&null_map);
+ ASSERT_EQ(2, null_map.size());
+ EXPECT_EQ(0, null_map[0]);
+ EXPECT_EQ(1, null_map[1]);
+}
+
TEST_F(VOrcTransformerTest, PreservesNullableArrayStructChildPositions) {
const auto int_type = make_nullable(std::make_shared<DataTypeInt32>());
const auto string_type = make_nullable(std::make_shared<DataTypeString>());
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]