This is an automated email from the ASF dual-hosted git repository.
Gabriel39 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 702b16cee25 [fix](iceberg) Handle NaN and -0.0 in float/double
predicate pushdown and identity partitions (#68128)
702b16cee25 is described below
commit 702b16cee256c20fcf79c4159cc35d140fa4790f
Author: daidai <[email protected]>
AuthorDate: Tue Sep 29 18:12:24 2026 +0800
[fix](iceberg) Handle NaN and -0.0 in float/double predicate pushdown and
identity partitions (#68128)
### What problem does this PR solve?
Issue Number: None
Related PR: None
Problem Summary:
Doris compares floats with "NaN is greater than everything, `NaN = NaN`"
plus IEEE zero equality (`-0.0 == 0.0`). Iceberg prunes files with
`Double.compare`: NaN is not in the bounds at all (the spec keeps it
only in `nan_value_counts`) and `-0.0` sorts strictly before `+0.0`. FE
translated predicates across that gap literally, which causes two bugs.
**1. Silently missing rows.** Iceberg prunes files that do hold matching
rows, and a pruned file never becomes a split, so BE cannot filter them
back in. With one spark-written data file per table:
| table | rows | query | result |
|---|---|---|---|
| `nan_only` | `(7, NaN)` | `where d > 0` / `d >= 0` | empty,
`inputSplitNum=0` |
| `nan_mixed` | `(1, 1.0), (2, NaN)` | `where d > 5` | empty — NaN hides
outside the bounds `[1.0, 1.0]` |
| `negzero` | `(7, -0.0)` | `where d = 0` / `d >= 0` | empty |
`select id, d, d > 0 from nan_only` returns `7, NaN, true`, so the
projection and the filter disagree.
The fix makes the shared FLOAT/DOUBLE leaves of
`IcebergPredicateConverter` emit the expression *equivalent* to the
Doris predicate under Iceberg's order. SCAN, CONFLICT and REWRITE mode
all use those leaves.
- `>` `>=` `!=` `NOT IN` gain `OR isNaN(col)`; `<` `<=` `=` `IN` gain
`AND notNaN(col)`. Both halves are needed — De Morgan maps one onto the
other, so `NOT` keeps working. With only the first, `RewriteNot` lowers
`NOT (d < 5)` to a bare `d >= 5` and prunes the NaN file again.
- A zero literal is bounded at the signed zero that reproduces IEEE:
`-0.0` for `>=` and `<`, `+0.0` for `>` and `<=`; `=` and `IN` expand to
both zeros. Unrelated files are still pruned.
- A NaN literal becomes `isNaN` / `notNaN`. `Expressions.*` rejects a
NaN literal outright, so `WHERE d > cast('nan' as double)` previously
failed during planning.
**2. `Cannot create expression literal from NaN` at commit.**
`IcebergPartitionUtils.parsePartitionValueFromString` maps BE's `nan`
spelling to `Double.NaN`, and both identity-partition predicate builders
passed it to `Expressions.equal`. So DELETE/UPDATE/MERGE touching a NaN
partition, and `INSERT OVERWRITE ... PARTITION(d='nan')`, planned fine
and then failed at commit — after the data files were written. Both now
share one helper that emits `isNaN` for NaN, `IS NULL` for null, and the
unchanged equality otherwise.
### Release note
Fixed Iceberg queries silently returning too few rows when a
FLOAT/DOUBLE column contains NaN or `-0.0`, fixed a planning failure
when the predicate literal is NaN, and fixed `Cannot create expression
literal from NaN` when writing to a NaN value in a FLOAT/DOUBLE identity
partition.
### Check List (For Author)
- Test
- [x] Regression test
- [x] Unit Test
- [ ] Manual test (add detailed scripts or steps below)
- [ ] No need to test or manual test. Explain why:
- Behavior changed:
- [ ] No.
- [x] Yes.
- Does this need documentation?
- [x] No.
- [ ] Yes.
---
be/src/exec/sink/viceberg_merge_sink.cpp | 3 +
.../writer/iceberg/viceberg_partition_writer.cpp | 7 +-
.../writer/iceberg/viceberg_partition_writer.h | 3 +
be/src/format/table/iceberg/nan_value_counter.cpp | 89 ++++
be/src/format/table/iceberg/nan_value_counter.h | 73 ++++
.../format/transformer/viceberg_parquet_writer.cpp | 26 +-
.../format/transformer/viceberg_parquet_writer.h | 16 +-
be/src/format/transformer/vorc_transformer.cpp | 18 +-
be/src/format/transformer/vorc_transformer.h | 9 +-
.../create_preinstalled_scripts/iceberg/run32.sql | 37 ++
.../iceberg/IcebergConnectorTransaction.java | 32 +-
.../iceberg/IcebergPredicateConverter.java | 228 +++++++++--
.../iceberg/IcebergWritePlanProvider.java | 8 +
.../connector/iceberg/IcebergWriterHelper.java | 50 ++-
.../iceberg/IcebergConnectorTransactionTest.java | 118 ++++++
...cebergPredicateConverterFloatSemanticsTest.java | 456 +++++++++++++++++++++
.../iceberg/IcebergWritePlanProviderTest.java | 24 ++
.../connector/iceberg/IcebergWriterHelperTest.java | 100 +++++
gensrc/thrift/DataSinks.thrift | 10 +
.../iceberg/write/test_iceberg_write_stats2.out | 4 +-
.../test_iceberg_float_predicate_pushdown.groovy | 124 ++++++
.../test_iceberg_write_nan_value_counts.groovy | 109 +++++
22 files changed, 1498 insertions(+), 46 deletions(-)
diff --git a/be/src/exec/sink/viceberg_merge_sink.cpp
b/be/src/exec/sink/viceberg_merge_sink.cpp
index 808bc5f1876..b025ce2cb83 100644
--- a/be/src/exec/sink/viceberg_merge_sink.cpp
+++ b/be/src/exec/sink/viceberg_merge_sink.cpp
@@ -410,6 +410,9 @@ Status VIcebergMergeSink::_build_inner_sinks() {
if (merge_sink.__isset.collect_column_stats) {
table_sink.__set_collect_column_stats(merge_sink.collect_column_stats);
}
+ if (merge_sink.__isset.nan_count_field_ids) {
+ table_sink.__set_nan_count_field_ids(merge_sink.nan_count_field_ids);
+ }
_table_sink.__set_type(TDataSinkType::ICEBERG_TABLE_SINK);
_table_sink.__set_iceberg_table_sink(table_sink);
diff --git a/be/src/exec/sink/writer/iceberg/viceberg_partition_writer.cpp
b/be/src/exec/sink/writer/iceberg/viceberg_partition_writer.cpp
index dfa2453c6aa..abdde43e637 100644
--- a/be/src/exec/sink/writer/iceberg/viceberg_partition_writer.cpp
+++ b/be/src/exec/sink/writer/iceberg/viceberg_partition_writer.cpp
@@ -53,6 +53,9 @@ VIcebergPartitionWriter::VIcebergPartitionWriter(
if (t_sink.iceberg_table_sink.__isset.collect_column_stats) {
_collect_column_stats = t_sink.iceberg_table_sink.collect_column_stats;
}
+ if (t_sink.iceberg_table_sink.__isset.nan_count_field_ids) {
+ _nan_count_field_ids = t_sink.iceberg_table_sink.nan_count_field_ids;
+ }
}
Status VIcebergPartitionWriter::open(RuntimeState* state, RuntimeProfile*
profile,
@@ -107,14 +110,14 @@ Status VIcebergPartitionWriter::open(RuntimeState* state,
RuntimeProfile* profil
.enable_int96_timestamps =
false};
_file_format_transformer = std::make_unique<VIcebergParquetWriter>(
state, _file_writer.get(), _write_output_expr_ctxs,
_write_column_names, false,
- parquet_options, _iceberg_schema_json, _schema);
+ parquet_options, _iceberg_schema_json, _schema,
_nan_count_field_ids);
open_status = _file_format_transformer->open();
break;
}
case TFileFormatType::FORMAT_ORC: {
_file_format_transformer = std::make_unique<VOrcTransformer>(
state, _file_writer.get(), _write_output_expr_ctxs, "",
_write_column_names, false,
- _compress_type, &_schema, _fs);
+ _compress_type, &_schema, _fs, _nan_count_field_ids);
open_status = _file_format_transformer->open();
break;
}
diff --git a/be/src/exec/sink/writer/iceberg/viceberg_partition_writer.h
b/be/src/exec/sink/writer/iceberg/viceberg_partition_writer.h
index f3a7a6002f8..954d458538c 100644
--- a/be/src/exec/sink/writer/iceberg/viceberg_partition_writer.h
+++ b/be/src/exec/sink/writer/iceberg/viceberg_partition_writer.h
@@ -101,6 +101,9 @@ private:
const std::map<std::string, std::string>& _hadoop_conf;
ClosedFileCallback _closed_file_callback;
bool _collect_column_stats = true;
+ // FE's metrics policy for NaN counting; empty (including an older FE that
sends nothing) means the
+ // parquet writer skips the extra pass entirely. See
IcebergWriterHelper#nanCountFieldIds.
+ std::vector<int32_t> _nan_count_field_ids;
std::shared_ptr<io::FileSystem> _fs = nullptr;
diff --git a/be/src/format/table/iceberg/nan_value_counter.cpp
b/be/src/format/table/iceberg/nan_value_counter.cpp
new file mode 100644
index 00000000000..d7b3efff4f0
--- /dev/null
+++ b/be/src/format/table/iceberg/nan_value_counter.cpp
@@ -0,0 +1,89 @@
+// 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 "format/table/iceberg/nan_value_counter.h"
+
+#include <cmath>
+#include <unordered_set>
+
+#include "common/check.h"
+#include "core/block/block.h"
+#include "core/column/column.h"
+#include "core/column/column_nullable.h"
+#include "core/column/column_vector.h"
+#include "format/table/iceberg/types.h"
+
+namespace doris::iceberg {
+
+namespace {
+
+// Branchless so the compiler can vectorize it: this runs over every floating
value written.
+template <typename Container>
+int64_t count_nan_values(const Container& data, const NullMap* null_map) {
+ int64_t nan_count = 0;
+ const size_t rows = data.size();
+ if (null_map == nullptr) {
+ for (size_t i = 0; i < rows; ++i) {
+ nan_count += static_cast<int64_t>(std::isnan(data[i]));
+ }
+ return nan_count;
+ }
+ for (size_t i = 0; i < rows; ++i) {
+ nan_count += static_cast<int64_t>((*null_map)[i] == 0 &&
std::isnan(data[i]));
+ }
+ return nan_count;
+}
+
+} // namespace
+
+NanValueCounter::NanValueCounter(const Schema& schema,
+ const std::vector<int32_t>&
requested_field_ids) {
+ const std::unordered_set<int32_t> requested(requested_field_ids.begin(),
+ requested_field_ids.end());
+ const auto& columns = schema.columns();
+ for (size_t i = 0; i < columns.size(); ++i) {
+ const auto type_id = columns[i].field_type()->type_id();
+ if (type_id != TypeID::FLOAT && type_id != TypeID::DOUBLE) {
+ continue;
+ }
+ const int32_t field_id = columns[i].field_id();
+ if (requested.contains(field_id)) {
+ _columns.emplace_back(i, field_id);
+ _counts[field_id] = 0;
+ }
+ }
+}
+
+void NanValueCounter::count(const Block& block) {
+ for (const auto& [column_position, field_id] : _columns) {
+ const IColumn* column =
block.get_by_position(column_position).column.get();
+ const NullMap* null_map = nullptr;
+ if (const auto* nullable =
check_and_get_column<ColumnNullable>(column)) {
+ null_map = &nullable->get_null_map_data();
+ column = &nullable->get_nested_column();
+ }
+ if (const auto* float64 = check_and_get_column<ColumnFloat64>(column))
{
+ _counts[field_id] += count_nan_values(float64->get_data(),
null_map);
+ continue;
+ }
+ const auto* float32 = check_and_get_column<ColumnFloat32>(column);
+ DORIS_CHECK(float32 != nullptr);
+ _counts[field_id] += count_nan_values(float32->get_data(), null_map);
+ }
+}
+
+} // namespace doris::iceberg
diff --git a/be/src/format/table/iceberg/nan_value_counter.h
b/be/src/format/table/iceberg/nan_value_counter.h
new file mode 100644
index 00000000000..5f7e35617a8
--- /dev/null
+++ b/be/src/format/table/iceberg/nan_value_counter.h
@@ -0,0 +1,73 @@
+// 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.
+
+#pragma once
+
+#include <cstdint>
+#include <map>
+#include <utility>
+#include <vector>
+
+#include "format/table/iceberg/schema.h"
+
+namespace doris {
+class Block;
+}
+
+namespace doris::iceberg {
+
+/**
+ * Counts NaNs per iceberg field as the rows go past, for one output file.
+ *
+ * Iceberg excludes NaN from a column's bounds by spec ("NaNs are not
permitted as lower or upper
+ * bounds"), so nan_value_counts is the ONLY metadata that can prove a file
holds no NaN. Without it
+ * InclusiveMetricsEvaluator.isNaN must assume NaN may be present and a float
range predicate cannot
+ * prune the file at all (see the FLOAT/DOUBLE leaves in FE
IcebergPredicateConverter). NEITHER parquet
+ * nor ORC column statistics carry a NaN count, so unlike every other metric
the writers report this one
+ * cannot be read back from the footer -- hence this counter, shared by both
writers so the semantics
+ * below exist once.
+ *
+ * A field this counter reports is a claim that every one of its values was
counted, which is what makes
+ * a reported zero trustworthy. Two independent narrowings apply, and a field
excluded by either one is
+ * simply absent from the result -- "unknown", which iceberg reads
conservatively -- rather than wrongly
+ * claimed NaN-free:
+ * - policy: {@code requested_field_ids}, the fields whose count FE would
keep under the table's
+ * metrics config. Counting one FE drops is a pure waste of a data pass,
and iceberg disables metrics
+ * for everything past the first 100 fields by default, so a wide table
hits this with no property
+ * set. An empty list (including an older FE that sends nothing) counts
nothing at all.
+ * - capability: only top-level columns, because a floating field nested in
a struct/list/map is not a
+ * block column of its own.
+ */
+class NanValueCounter {
+public:
+ NanValueCounter(const Schema& schema, const std::vector<int32_t>&
requested_field_ids);
+
+ // Block column i maps to iceberg column i: both writers build their
output schema straight from
+ // Schema::columns() in order, and both reject a column-count mismatch
before reaching here.
+ void count(const Block& block);
+
+ const std::map<int, int64_t>& counts() const { return _counts; }
+
+ bool empty() const { return _counts.empty(); }
+
+private:
+ // (block column position, iceberg field id) of the columns counted for
this file.
+ std::vector<std::pair<size_t, int32_t>> _columns;
+ std::map<int, int64_t> _counts;
+};
+
+} // namespace doris::iceberg
diff --git a/be/src/format/transformer/viceberg_parquet_writer.cpp
b/be/src/format/transformer/viceberg_parquet_writer.cpp
index 57f329014d0..ceff37500a5 100644
--- a/be/src/format/transformer/viceberg_parquet_writer.cpp
+++ b/be/src/format/transformer/viceberg_parquet_writer.cpp
@@ -51,18 +51,33 @@ static std::string encode_iceberg_bound(const
::parquet::Statistics& stats, std:
encoded.erase(0, start);
return encoded;
}
-
VIcebergParquetWriter::VIcebergParquetWriter(RuntimeState* state,
io::FileWriter* file_writer,
const VExprContextSPtrs&
output_vexpr_ctxs,
std::vector<std::string>
column_names,
bool output_object_data,
const ParquetFileOptions&
parquet_options,
const std::string*
iceberg_schema_json,
- const iceberg::Schema&
iceberg_schema)
+ const iceberg::Schema&
iceberg_schema,
+ const std::vector<int32_t>&
nan_count_field_ids)
: VParquetWriter(state, file_writer, output_vexpr_ctxs,
std::move(column_names),
output_object_data, parquet_options),
_iceberg_schema(iceberg_schema),
- _iceberg_schema_json(iceberg_schema_json == nullptr ? "" :
*iceberg_schema_json) {}
+ _iceberg_schema_json(iceberg_schema_json == nullptr ? "" :
*iceberg_schema_json),
+ _nan_count_field_ids(nan_count_field_ids) {}
+
+Status VIcebergParquetWriter::open() {
+ RETURN_IF_ERROR(VParquetWriter::open());
+ _nan_value_counter.emplace(_iceberg_schema, _nan_count_field_ids);
+ return Status::OK();
+}
+
+Status VIcebergParquetWriter::write(const Block& block) {
+ if (block.rows() == 0) {
+ return Status::OK();
+ }
+ _nan_value_counter->count(block);
+ return VParquetWriter::write(block);
+}
std::unique_ptr<ArrowBlockConvertor>
VIcebergParquetWriter::_create_arrow_block_convertor(
DataTypes types, std::vector<std::string> names, const std::string&
timezone_name,
@@ -123,6 +138,11 @@ Status
VIcebergParquetWriter::collect_file_statistics_after_close(TIcebergColumn
stats->__set_column_sizes(column_sizes);
stats->__set_value_counts(value_counts);
+ // Left unset when no column was counted, so FE keeps reporting "unknown"
rather than an empty
+ // claim -- the same shape an older BE produces.
+ if (!_nan_value_counter->empty()) {
+ stats->__set_nan_value_counts(_nan_value_counter->counts());
+ }
if (has_any_null_count) {
stats->__set_null_value_counts(null_value_counts);
}
diff --git a/be/src/format/transformer/viceberg_parquet_writer.h
b/be/src/format/transformer/viceberg_parquet_writer.h
index 4905d49510b..787585eb569 100644
--- a/be/src/format/transformer/viceberg_parquet_writer.h
+++ b/be/src/format/transformer/viceberg_parquet_writer.h
@@ -17,7 +17,12 @@
#pragma once
+#include <cstdint>
+#include <optional>
+#include <vector>
+
#include "format/table/iceberg/iceberg_arrow_block_convertor.h"
+#include "format/table/iceberg/nan_value_counter.h"
#include "format/table/iceberg/schema.h"
#include "format/transformer/vparquet_writer.h"
@@ -30,7 +35,12 @@ public:
std::vector<std::string> column_names, bool
output_object_data,
const ParquetFileOptions& parquet_options,
const std::string* iceberg_schema_json,
- const iceberg::Schema& iceberg_schema);
+ const iceberg::Schema& iceberg_schema,
+ const std::vector<int32_t>& nan_count_field_ids =
{});
+
+ Status open() override;
+
+ Status write(const Block& block) override;
Status collect_file_statistics_after_close(TIcebergColumnStats* stats);
@@ -42,6 +52,10 @@ protected:
private:
const iceberg::Schema& _iceberg_schema;
std::string _iceberg_schema_json;
+ const std::vector<int32_t> _nan_count_field_ids;
+ // Built at open(), once the writer is past schema validation. See
iceberg::NanValueCounter for why a
+ // reported zero is a claim and which fields are counted at all.
+ std::optional<iceberg::NanValueCounter> _nan_value_counter;
};
} // namespace doris
diff --git a/be/src/format/transformer/vorc_transformer.cpp
b/be/src/format/transformer/vorc_transformer.cpp
index 8f1db599fc9..992d5a040e8 100644
--- a/be/src/format/transformer/vorc_transformer.cpp
+++ b/be/src/format/transformer/vorc_transformer.cpp
@@ -182,14 +182,16 @@ VOrcTransformer::VOrcTransformer(RuntimeState* state,
doris::io::FileWriter* fil
std::vector<std::string> column_names, bool
output_object_data,
TFileCompressType::type compress_type,
const iceberg::Schema* iceberg_schema,
- std::shared_ptr<io::FileSystem> fs)
+ std::shared_ptr<io::FileSystem> fs,
+ const std::vector<int32_t>&
nan_count_field_ids)
: VFileFormatTransformer(state, output_vexpr_ctxs, output_object_data),
_fs(fs),
_file_writer(file_writer),
_column_names(std::move(column_names)),
_write_options(new orc::WriterOptions()),
_schema_str(std::move(schema)),
- _iceberg_schema(iceberg_schema) {
+ _iceberg_schema(iceberg_schema),
+ _nan_count_field_ids(nan_count_field_ids) {
// ORC recognizes GMT as its UTC fast path. Other UTC aliases can lack a
zoneinfo file
// or resolve through a locally overridden UTC file, shifting timestamp
statistics.
int32_t fixed_offset = 0;
@@ -249,6 +251,9 @@ Status VOrcTransformer::open() {
return Status::InternalError("failed to create writer: {}", e.what());
}
_writer->addUserMetadata("CreatedBy", doris::get_short_version());
+ if (_iceberg_schema != nullptr) {
+ _nan_value_counter.emplace(*_iceberg_schema, _nan_count_field_ids);
+ }
return Status::OK();
}
@@ -537,6 +542,12 @@ Status
VOrcTransformer::collect_file_statistics_after_close(TIcebergColumnStats*
}
stats->__set_value_counts(value_counts);
+ // ORC statistics carry no NaN count, so it comes from the counter fed
during write() rather than
+ // from the footer read above. Left unset when no column was counted,
so FE keeps reporting
+ // "unknown" instead of an empty claim -- the same shape an older BE
produces.
+ if (_nan_value_counter.has_value() && !_nan_value_counter->empty()) {
+ stats->__set_nan_value_counts(_nan_value_counter->counts());
+ }
if (has_any_null_count) {
stats->__set_null_value_counts(null_value_counts);
}
@@ -694,6 +705,9 @@ Status VOrcTransformer::write(const Block& block) {
if (block.rows() == 0) {
return Status::OK();
}
+ if (_nan_value_counter.has_value()) {
+ _nan_value_counter->count(block);
+ }
// Buffer used by date/datetime/datev2/datetimev2/largeint type
Arena arena;
diff --git a/be/src/format/transformer/vorc_transformer.h
b/be/src/format/transformer/vorc_transformer.h
index 2653f4c6019..6515439cc2c 100644
--- a/be/src/format/transformer/vorc_transformer.h
+++ b/be/src/format/transformer/vorc_transformer.h
@@ -21,6 +21,7 @@
#include <stdint.h>
#include <memory>
+#include <optional>
#include <orc/OrcFile.hh>
#include <string>
#include <vector>
@@ -28,6 +29,7 @@
#include "common/status.h"
#include "core/block/block.h"
#include "core/column/column_nullable.h"
+#include "format/table/iceberg/nan_value_counter.h"
#include "format/table/iceberg/schema.h"
#include "format/transformer/vparquet_writer.h"
#include "orc/Type.hh"
@@ -85,7 +87,8 @@ public:
std::vector<std::string> column_names, bool
output_object_data,
TFileCompressType::type compression,
const iceberg::Schema* iceberg_schema = nullptr,
- std::shared_ptr<io::FileSystem> fs = nullptr);
+ std::shared_ptr<io::FileSystem> fs = nullptr,
+ const std::vector<int32_t>& nan_count_field_ids = {});
~VOrcTransformer() = default;
@@ -126,6 +129,10 @@ private:
const iceberg::Schema* _iceberg_schema;
std::vector<uint8_t> _iceberg_binary_normalization_required;
+ // ORC column statistics carry no NaN count either, so it is accumulated
while the rows go past --
+ // see iceberg::NanValueCounter. Only populated for an iceberg write.
+ const std::vector<int32_t> _nan_count_field_ids;
+ std::optional<iceberg::NanValueCounter> _nan_value_counter;
// Buffer used by date/datetime/datev2/datetimev2/largeint type
// date/datetime/datev2/datetimev2/largeint type will be converted to
string bytes to store in Buffer
diff --git
a/docker/thirdparties/docker-compose/iceberg/scripts/create_preinstalled_scripts/iceberg/run32.sql
b/docker/thirdparties/docker-compose/iceberg/scripts/create_preinstalled_scripts/iceberg/run32.sql
new file mode 100644
index 00000000000..64506ce0776
--- /dev/null
+++
b/docker/thirdparties/docker-compose/iceberg/scripts/create_preinstalled_scripts/iceberg/run32.sql
@@ -0,0 +1,37 @@
+create database if not exists demo.test_db;
+use demo.test_db;
+
+-- Fixtures for float/double predicate pushdown.
+--
+-- Iceberg keeps NaN out of a column's bounds entirely (spec: "NaNs are not
permitted as lower or upper
+-- bounds" -- it is recorded only in nan_value_counts) and orders -0.0
strictly before +0.0, while Doris
+-- matches rows with "NaN is greater than everything" and IEEE zero equality
(-0.0 == 0.0). A predicate
+-- translated literally therefore prunes files that do hold matching rows.
+--
+-- COALESCE(1) keeps every table at exactly one data file, so a wrongly pruned
file shows up directly as
+-- inputSplitNum=0 in EXPLAIN. These must be written by spark, not by Doris:
Doris reports no
+-- nan_value_counts at all, and its parquet writer normalizes zero bounds, so
Doris-written files do not
+-- carry the metadata shape that triggers the pruning.
+
+drop table if exists float_prune_nan_only;
+create table float_prune_nan_only (id int, d double) using iceberg;
+insert into float_prune_nan_only
+select /*+ COALESCE(1) */ * from values (7, cast('NaN' as double)) as v(id, d);
+
+-- NaN hides outside the bounds, so this file is pruned by the bounds alone
for `d > 5`.
+drop table if exists float_prune_nan_mixed;
+create table float_prune_nan_mixed (id int, d double) using iceberg;
+insert into float_prune_nan_mixed
+select /*+ COALESCE(1) */ * from values (1, cast(1.0 as double)), (2,
cast('NaN' as double)) as v(id, d);
+
+drop table if exists float_prune_nan_mixed_float;
+create table float_prune_nan_mixed_float (id int, f float) using iceberg;
+insert into float_prune_nan_mixed_float
+select /*+ COALESCE(1) */ * from values (1, cast(1.0 as float)), (2,
cast('NaN' as float)) as v(id, f);
+
+-- Iceberg keeps the sign in the bounds (lower = upper = -0.0), so a `d = 0` /
`d >= 0` bound placed at +0.0
+-- sorts strictly after it and prunes the file.
+drop table if exists float_prune_negzero;
+create table float_prune_negzero (id int, d double) using iceberg;
+insert into float_prune_negzero
+select /*+ COALESCE(1) */ * from values (7, cast('-0.0' as double)) as v(id,
d);
diff --git
a/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergConnectorTransaction.java
b/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergConnectorTransaction.java
index a94156a3cf8..bd85c627e75 100644
---
a/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergConnectorTransaction.java
+++
b/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergConnectorTransaction.java
@@ -902,11 +902,7 @@ public class IcebergConnectorTransaction implements
ConnectorTransaction, Rewrit
}
Object partitionValue =
IcebergPartitionUtils.parsePartitionValueFromString(
partitionValueStr, sourceField.type(), zone);
- String sourceColName = sourceField.name();
- Expression eqExpr = partitionValue == null
- ? Expressions.isNull(sourceColName)
- : Expressions.equal(sourceColName, partitionValue);
- predicates.add(eqExpr);
+ predicates.add(identityPartitionPredicate(sourceField.name(),
partitionValue));
}
}
@@ -922,6 +918,28 @@ public class IcebergConnectorTransaction implements
ConnectorTransaction, Rewrit
return result;
}
+ /**
+ * {@code sourceCol = value} for one identity partition key, or the unary
predicate the value demands.
+ *
+ * <p>A NaN cannot be an iceberg literal at all — {@code Literals.from}
throws "Cannot create expression
+ * literal from NaN", and iceberg models it only through {@code
isNaN}/{@code notNaN}. Meanwhile
+ * {@link IcebergPartitionUtils#parsePartitionValueFromString}
deliberately parses Doris's {@code nan}
+ * spelling into {@link Double#NaN}, so a FLOAT/DOUBLE identity partition
holding NaN used to abort the
+ * whole commit ("Failed to commit iceberg transaction: Cannot create
expression literal from NaN") on
+ * both paths that build this predicate: DELETE/UPDATE/MERGE conflict
detection and
+ * {@code INSERT OVERWRITE ... PARTITION(d='nan')}.
+ */
+ private static Expression identityPartitionPredicate(String sourceColName,
Object value) {
+ if (value == null) {
+ return Expressions.isNull(sourceColName);
+ }
+ if ((value instanceof Double || value instanceof Float)
+ && Double.isNaN(((Number) value).doubleValue())) {
+ return Expressions.isNaN(sourceColName);
+ }
+ return Expressions.equal(sourceColName, value);
+ }
+
/**
* DELETE: commit position-delete files via {@link RowDelta}, guarded by
the optimistic conflict-detection
* validation suite ({@link #applyRowDeltaValidations}) and — on a V3
table — the deletion-vector
@@ -1205,9 +1223,7 @@ public class IcebergConnectorTransaction implements
ConnectorTransaction, Rewrit
valueStr = null;
}
Object value =
IcebergPartitionUtils.parsePartitionValueFromString(valueStr,
sourceField.type(), zone);
- Expression predicate = value == null
- ? Expressions.isNull(sourceField.name())
- : Expressions.equal(sourceField.name(), value);
+ Expression predicate =
identityPartitionPredicate(sourceField.name(), value);
expression = expression == null ? predicate :
Expressions.and(expression, predicate);
}
return expression;
diff --git
a/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergPredicateConverter.java
b/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergPredicateConverter.java
index 329023ead0b..7fd94a43125 100644
---
a/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergPredicateConverter.java
+++
b/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergPredicateConverter.java
@@ -48,6 +48,8 @@ import java.time.LocalDateTime;
import java.time.ZoneId;
import java.time.ZoneOffset;
import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Locale;
@@ -83,6 +85,10 @@ import java.util.Set;
* re-filter, so any unrepresentable sub-node collapses the whole expression
to {@code null} (the rewrite planner
* then fails loud rather than silently widening the set of files rewritten).
Unlike conflict mode it keeps
* cross-column {@code OR} / {@code NOT(comparison)} / {@code NE} and drops
the structural/UUID narrowing.</p>
+ *
+ * <p>All three modes share the FLOAT/DOUBLE leaves ({@link
#buildLeafComparison} / {@link #buildLeafIn}),
+ * which reconcile Doris's row-level float semantics with the total order
iceberg prunes files by. See the
+ * banner comment on those methods — getting this wrong silently drops rows,
not just performance.</p>
*/
public class IcebergPredicateConverter {
@@ -92,6 +98,10 @@ public class IcebergPredicateConverter {
private static final String ICEBERG_ROW_ID_COL = "_row_id";
private static final String ICEBERG_LAST_UPDATED_SEQUENCE_NUMBER_COL =
"_last_updated_sequence_number";
+ // One value to Doris (IEEE: -0.0 == 0.0), two adjacent points in
iceberg's total order. See the FLOAT/DOUBLE
+ // banner on buildLeafComparison.
+ private static final List<Object> SIGNED_ZEROS =
Collections.unmodifiableList(Arrays.asList(-0.0d, 0.0d));
+
private final Schema schema;
private final ZoneId sessionZone;
// The conversion matrix. SCAN = scan-time pushdown (BE re-filters, so
widening is safe). CONFLICT = O5-2
@@ -249,23 +259,7 @@ public class IcebergPredicateConverter {
}
return null;
}
- switch (cmp.getOperator()) {
- case EQ:
- case EQ_FOR_NULL:
- return Expressions.equal(colName, value);
- case NE:
- return Expressions.not(Expressions.equal(colName, value));
- case GE:
- return Expressions.greaterThanOrEqual(colName, value);
- case GT:
- return Expressions.greaterThan(colName, value);
- case LE:
- return Expressions.lessThanOrEqual(colName, value);
- case LT:
- return Expressions.lessThan(colName, value);
- default:
- return null;
- }
+ return buildLeafComparison(field.type(), colName, cmp.getOperator(),
value);
}
private Expression buildIn(ConnectorIn in) {
@@ -289,7 +283,191 @@ public class IcebergPredicateConverter {
}
values.add(value);
}
- return in.isNegated() ? Expressions.notIn(colName, values) :
Expressions.in(colName, values);
+ return buildLeafIn(field.type(), colName, values, in.isNegated());
+ }
+
+ //
════════════════════════════════════════════════════════════════════════════════════════════════════
+ // FLOAT/DOUBLE leaves: Doris row semantics vs the total order iceberg
prunes files by.
+ //
+ // Doris evaluates a ROW with "NaN is greater than everything, NaN = NaN"
plus IEEE zero equality
+ // (-0.0 == 0.0). Iceberg prunes a FILE with Comparators.naturalOrder()
(Double.compare), where NaN is
+ // not in the bounds at all -- the spec says "NaNs are not permitted as
lower or upper bounds", they are
+ // recorded separately in nan_value_counts -- and -0.0 sorts strictly
before +0.0. Handing iceberg a
+ // literal translation therefore prunes files that do hold matching rows,
and a pruned file never becomes
+ // a split, so BE's residual filter cannot recover those rows: the query
silently returns too few rows.
+ //
+ // The leaves below do not merely widen, they emit the expression
EQUIVALENT to the Doris predicate under
+ // iceberg's order, and they spell out the NaN half on BOTH sides: a
comparison Doris matches NaN with gets
+ // `OR isNaN`, one it does not gets `AND notNaN`. That symmetry is
load-bearing, not decoration -- De Morgan
+ // maps each form onto the other, so `NOT` composes at any nesting depth.
Leaving the `AND notNaN` off would
+ // still read correctly row by row, yet iceberg's RewriteNot would lower
not(lessThan(d, 5)) to a BARE
+ // gtEq(d, 5) whose file-level evaluator prunes a {1.0, NaN} file -- the
very bug, back through the negation.
+ // notNaN costs no pruning either: it only rules out a file that is
entirely NaN.
+ //
+ // Cost: InclusiveMetricsEvaluator.isNaN only prunes a file that reports
nan_value_count == 0, so files
+ // whose metrics omit NaN counts -- including every file Doris writes
today, IcebergWriterHelper passes a
+ // null nanValueCounts -- stop being pruned by a float range predicate.
This is the trade parquet-java and
+ // parquet-cpp already make ("stats may hold NaN -> do not prune");
writers that do report NaN counts
+ // (spark, flink, iceberg-java) keep pruning in full.
+ //
════════════════════════════════════════════════════════════════════════════════════════════════════
+
+ /**
+ * {@code col OP value}, shared by all three modes. A non-floating column
gets the plain 1:1 mapping; a
+ * FLOAT/DOUBLE column gets the NaN / signed-zero reconciliation described
above.
+ */
+ private static Expression buildLeafComparison(Type type, String colName,
+ ConnectorComparison.Operator op, Object value) {
+ if (!isFloating(type)) {
+ return buildPlainComparison(colName, op, value);
+ }
+ // On a floating column extractIcebergLiteral only ever yields a
Number: Double/Float from a float
+ // literal, Double from a decimal literal, Integer/Long from an
integer one (`d > 0`).
+ double v = ((Number) value).doubleValue();
+ if (Double.isNaN(v)) {
+ // Expressions.*(col, NaN) throws ("Cannot create expression
literal from NaN") -- iceberg models a
+ // NaN literal only through the unary isNaN/notNaN. Doris: NaN is
the greatest value, NaN = NaN.
+ //
+ // notNaN is the one arm that needs a NULL guard: iceberg reads
notNaN(null) as TRUE
+ // (Evaluator: !NaNUtil.isNaN(value)), while Doris leaves `NULL !=
NaN` / `NULL < NaN` UNKNOWN,
+ // so a null row does not match. isNaN and notNull are already
false for null. Without the guard
+ // a REWRITE, whose filter has no downstream re-filter, would pull
in a null-only file that holds
+ // no matching row at all.
+ switch (op) {
+ case EQ:
+ case EQ_FOR_NULL:
+ case GE:
+ return Expressions.isNaN(colName);
+ case NE:
+ case LT:
+ return andNotNull(colName, Expressions.notNaN(colName));
+ case GT:
+ return Expressions.alwaysFalse();
+ case LE:
+ return Expressions.notNull(colName);
+ default:
+ return null;
+ }
+ }
+ switch (op) {
+ case EQ:
+ case EQ_FOR_NULL:
+ return andNotNaN(colName, v == 0.0d
+ ? Expressions.in(colName, SIGNED_ZEROS) :
Expressions.equal(colName, value));
+ case NE:
+ return orIsNaN(colName, v == 0.0d
+ ? Expressions.notIn(colName, SIGNED_ZEROS)
+ : Expressions.not(Expressions.equal(colName, value)));
+ case GE:
+ // `d >= ±0.0` matches -0.0 too -> bound at -0.0. NaN
satisfies >= every literal.
+ return orIsNaN(colName,
Expressions.greaterThanOrEqual(colName, zeroBound(v, value, -0.0d)));
+ case GT:
+ // `d > ±0.0` matches neither zero -> bound at +0.0. NaN
satisfies > every literal.
+ return orIsNaN(colName, Expressions.greaterThan(colName,
zeroBound(v, value, 0.0d)));
+ case LE:
+ // `d <= ±0.0` matches +0.0 too -> bound at +0.0.
+ return andNotNaN(colName, Expressions.lessThanOrEqual(colName,
zeroBound(v, value, 0.0d)));
+ case LT:
+ // `d < ±0.0` matches neither zero -> bound at -0.0.
+ return andNotNaN(colName, Expressions.lessThan(colName,
zeroBound(v, value, -0.0d)));
+ default:
+ return null;
+ }
+ }
+
+ /**
+ * {@code col IN (values)} / {@code col NOT IN (values)}, shared by all
three modes. On a FLOAT/DOUBLE
+ * column a zero element expands to both signed zeros, a NaN element is
lifted out into an isNaN/notNaN
+ * arm (iceberg rejects a NaN literal), and NOT IN additionally keeps
NaN-bearing files, because Doris
+ * reads `NaN != v` as true.
+ */
+ private static Expression buildLeafIn(Type type, String colName,
List<Object> values, boolean negated) {
+ if (!isFloating(type)) {
+ return negated ? Expressions.notIn(colName, values) :
Expressions.in(colName, values);
+ }
+ List<Object> literals = new ArrayList<>(values.size() + 1);
+ boolean hasNaN = false;
+ boolean hasZero = false;
+ for (Object value : values) {
+ double v = ((Number) value).doubleValue();
+ if (Double.isNaN(v)) {
+ hasNaN = true;
+ } else if (v == 0.0d) {
+ hasZero = true;
+ } else {
+ literals.add(value);
+ }
+ }
+ if (hasZero) {
+ literals.addAll(SIGNED_ZEROS);
+ }
+ if (negated) {
+ // A listed NaN excludes NaN rows; otherwise NaN rows satisfy NOT
IN. iceberg's and()/or() fold the
+ // alwaysTrue/alwaysFalse identity away, so `d NOT IN (NaN)` comes
out as and(notNaN, notNull) --
+ // the same NULL guard the NaN-literal comparison needs, since
notNaN(null) reads as TRUE in
+ // iceberg while Doris leaves `NULL NOT IN (NaN)` UNKNOWN.
+ Expression excluded = literals.isEmpty()
+ ? Expressions.alwaysTrue() : Expressions.notIn(colName,
literals);
+ return hasNaN
+ ? andNotNull(colName, Expressions.and(excluded,
Expressions.notNaN(colName)))
+ : orIsNaN(colName, excluded);
+ }
+ Expression included = literals.isEmpty()
+ ? Expressions.alwaysFalse() : andNotNaN(colName,
Expressions.in(colName, literals));
+ return hasNaN ? orIsNaN(colName, included) : included;
+ }
+
+ // The plain 1:1 mapping, i.e. what every mode emitted for every column
before the FLOAT/DOUBLE leaves.
+ private static Expression buildPlainComparison(String colName,
ConnectorComparison.Operator op,
+ Object value) {
+ switch (op) {
+ case EQ:
+ case EQ_FOR_NULL:
+ return Expressions.equal(colName, value);
+ case NE:
+ return Expressions.not(Expressions.equal(colName, value));
+ case GE:
+ return Expressions.greaterThanOrEqual(colName, value);
+ case GT:
+ return Expressions.greaterThan(colName, value);
+ case LE:
+ return Expressions.lessThanOrEqual(colName, value);
+ case LT:
+ return Expressions.lessThan(colName, value);
+ default:
+ return null;
+ }
+ }
+
+ private static boolean isFloating(Type type) {
+ TypeID id = type.typeId();
+ return id == TypeID.FLOAT || id == TypeID.DOUBLE;
+ }
+
+ // iceberg orders -0.0 strictly before +0.0, so bounding at the wrong zero
drops the other one. Pick the
+ // signed zero whose total-order bound covers exactly the IEEE matching
set; pass other literals through.
+ private static Object zeroBound(double v, Object value, double signedZero)
{
+ return v == 0.0d ? signedZero : value;
+ }
+
+ // Doris matches NaN rows for this comparison, but iceberg keeps NaN out
of the file bounds entirely --
+ // this arm is what stops a NaN-bearing file from being pruned.
+ private static Expression orIsNaN(String colName, Expression expr) {
+ return Expressions.or(expr, Expressions.isNaN(colName));
+ }
+
+ // Doris does NOT match NaN rows for this comparison. Saying so explicitly
is redundant going forward (the
+ // bounds already exclude NaN) but is exactly what makes the leaf
self-dual: De Morgan turns
+ // and(pred, notNaN) into or(!pred, isNaN), so a NOT wrapping this leaf
still keeps NaN-bearing files
+ // instead of collapsing to a bare range predicate. See the banner above.
+ private static Expression andNotNaN(String colName, Expression expr) {
+ return Expressions.and(expr, Expressions.notNaN(colName));
+ }
+
+ // Only needed where the emitted arm is a bare notNaN: iceberg evaluates
notNaN(null) as TRUE, while Doris
+ // leaves a comparison against NULL UNKNOWN, so the row does not match.
Every other leaf already excludes
+ // null on its own -- a range/equality bound is false for null, and
isNaN/notNull are false for null.
+ private static Expression andNotNull(String colName, Expression expr) {
+ return Expressions.and(expr, Expressions.notNull(colName));
}
private Types.NestedField getPushdownField(String colName) {
@@ -704,7 +882,7 @@ public class IcebergPredicateConverter {
if (isUuid(type) && !values.isEmpty()) {
return null;
}
- Expression valuesExpr = values.isEmpty() ? null :
Expressions.in(field.name(), values);
+ Expression valuesExpr = values.isEmpty() ? null : buildLeafIn(type,
field.name(), values, false);
Expression nullExpr = hasNull ? Expressions.isNull(field.name()) :
null;
return combineOr(nullExpr, valuesExpr);
}
@@ -732,7 +910,8 @@ public class IcebergPredicateConverter {
}
String colName = field.name();
return Expressions.and(
- Expressions.greaterThanOrEqual(colName, lo),
Expressions.lessThanOrEqual(colName, hi));
+ buildLeafComparison(type, colName,
ConnectorComparison.Operator.GE, lo),
+ buildLeafComparison(type, colName,
ConnectorComparison.Operator.LE, hi));
}
private Expression buildConflictComparison(ConnectorComparison cmp) {
@@ -774,15 +953,11 @@ public class IcebergPredicateConverter {
}
switch (cmp.getOperator()) {
case EQ:
- return Expressions.equal(colName, value);
case GT:
- return Expressions.greaterThan(colName, value);
case GE:
- return Expressions.greaterThanOrEqual(colName, value);
case LT:
- return Expressions.lessThan(colName, value);
case LE:
- return Expressions.lessThanOrEqual(colName, value);
+ return buildLeafComparison(type, colName, cmp.getOperator(),
value);
default:
// NE / EQ_FOR_NULL are not part of the legacy conflict matrix
-> dropped.
return null;
@@ -877,7 +1052,8 @@ public class IcebergPredicateConverter {
}
String colName = field.name();
return Expressions.and(
- Expressions.greaterThanOrEqual(colName, lo),
Expressions.lessThanOrEqual(colName, hi));
+ buildLeafComparison(field.type(), colName,
ConnectorComparison.Operator.GE, lo),
+ buildLeafComparison(field.type(), colName,
ConnectorComparison.Operator.LE, hi));
}
private static Expression combineOr(Expression left, Expression right) {
diff --git
a/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergWritePlanProvider.java
b/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergWritePlanProvider.java
index 2611f56b074..11a439b1a36 100644
---
a/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergWritePlanProvider.java
+++
b/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergWritePlanProvider.java
@@ -741,6 +741,10 @@ public class IcebergWritePlanProvider implements
ConnectorWritePlanProvider {
// table (all columns metrics=none) skips BE collection entirely. Same
schema the sink advertises.
tSink.setCollectColumnStats(
IcebergWriterHelper.shouldCollectColumnStats(schemaContext,
schemaContext.getSchema()));
+ // A NaN count is the one statistic the parquet footer does not carry,
so BE pays an extra pass for it.
+ // The table-wide flag above is too coarse to gate that pass -- see
IcebergWriterHelper#nanCountFieldIds.
+ tSink.setNanCountFieldIds(
+ IcebergWriterHelper.nanCountFieldIds(schemaContext,
schemaContext.getSchema()));
// Partition spec (only for a partitioned table, mirroring legacy
spec().isPartitioned()).
PartitionSpec partitionSpec = schemaContext.getPartitionSpec();
@@ -812,6 +816,7 @@ public class IcebergWritePlanProvider implements
ConnectorWritePlanProvider {
tSink.setSchemaJson(SchemaParser.toJson(rewriteSchema));
// #65782: the collect flag must reflect the same (v3-appended)
schema the sink advertises.
tSink.setCollectColumnStats(IcebergWriterHelper.shouldCollectColumnStats(table,
rewriteSchema));
+
tSink.setNanCountFieldIds(IcebergWriterHelper.nanCountFieldIds(table,
rewriteSchema));
}
return tSink;
}
@@ -883,6 +888,9 @@ public class IcebergWritePlanProvider implements
ConnectorWritePlanProvider {
// #65782: gate BE-side column-stats collection on the table's iceberg
metrics policy (v3-appended schema).
tSink.setCollectColumnStats(
IcebergWriterHelper.shouldCollectColumnStats(schemaContext,
schema));
+ // The replacement data files UPDATE / SQL MERGE write reach the same
iceberg parquet writer as an
+ // INSERT, so they need the same NaN-count policy (against the MERGE
schema) to stay prunable.
+
tSink.setNanCountFieldIds(IcebergWriterHelper.nanCountFieldIds(schemaContext,
schema));
// #66112: UPDATE and SQL MERGE share this sink, but only SQL MERGE
has the one-source-row invariant.
tSink.setRequireMergeCardinalityCheck(requireMergeCardinalityCheck);
tSink.setWritesDataFiles(writesDataFiles);
diff --git
a/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergWriterHelper.java
b/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergWriterHelper.java
index 8e0839c786f..a10a95cb45b 100644
---
a/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergWriterHelper.java
+++
b/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergWriterHelper.java
@@ -193,6 +193,11 @@ final class IcebergWriterHelper {
Map<Integer, Long> nullValueCounts = new HashMap<>();
Map<Integer, ByteBuffer> lowerBounds = new HashMap<>();
Map<Integer, ByteBuffer> upperBounds = new HashMap<>();
+ // Deliberately null rather than empty when BE reports nothing:
iceberg reads a missing NaN count as
+ // "may contain NaN" and keeps the file for a float range predicate,
which is the only safe reading of
+ // a BE that does not count NaNs (an older one, or a format whose
writer cannot). An explicit zero is
+ // a positive claim that the column holds none, so it must only ever
come from BE having counted.
+ Map<Integer, Long> nanValueCounts = null;
if (commitData.isSetColumnStats()) {
TIcebergColumnStats stats = commitData.column_stats;
if (stats.isSetColumnSizes()) {
@@ -204,6 +209,9 @@ final class IcebergWriterHelper {
if (stats.isSetNullValueCounts()) {
nullValueCounts = stats.null_value_counts;
}
+ if (stats.isSetNanValueCounts()) {
+ nanValueCounts = stats.nan_value_counts;
+ }
if (stats.isSetLowerBounds()) {
lowerBounds = stats.lower_bounds;
}
@@ -217,7 +225,8 @@ final class IcebergWriterHelper {
filterDisabledMetrics(columnSizes, schema, metricsConfig),
filterLogicalMetrics(valueCounts, schema, metricsConfig,
fieldParents),
filterLogicalMetrics(nullValueCounts, schema, metricsConfig,
fieldParents),
- null,
+ nanValueCounts == null ? null
+ : filterLogicalMetrics(nanValueCounts, schema,
metricsConfig, fieldParents),
filterBounds(lowerBounds, schema, metricsConfig, fieldParents,
fileFormat, true),
filterBounds(upperBounds, schema, metricsConfig, fieldParents,
fileFormat, false));
}
@@ -556,4 +565,43 @@ final class IcebergWriterHelper {
.anyMatch(field -> MetricsUtil.metricsMode(writerSchema,
metricsConfig, field.fieldId())
!= MetricsModes.None.get());
}
+
+ /**
+ * Field ids of the FLOAT/DOUBLE fields whose NaN count would survive this
table's metrics policy, i.e.
+ * whose effective mode is not {@code none}. Threaded to the BE so it
counts only what
+ * {@link #buildDataFileMetrics} would keep.
+ *
+ * <p>This exists because a NaN count is the one statistic BE cannot read
back from the parquet footer --
+ * it is an extra pass over the values. {@link #shouldCollectColumnStats}
is table-wide and stays true as
+ * soon as any single field is enabled, which is not a fine enough gate:
iceberg gives the first
+ * {@code write.metadata.metrics.max-inferred-column-defaults} (100)
fields the default mode and
+ * {@code none} to everything after it, so a wide table disables most of
its columns with no property set
+ * at all. Without this list BE would scan every float column on every
block and FE would then drop most
+ * of the results.</p>
+ *
+ * <p>The list is the metrics POLICY, not a capability claim: BE counts
the subset it can actually reach
+ * (today, top-level columns), and a field it does not count simply stays
absent from the manifest, which
+ * iceberg reads as "may contain NaN".</p>
+ */
+ static List<Integer> nanCountFieldIds(Table table, Schema writerSchema) {
+ return nanCountFieldIds(MetricsConfig.forTable(table), writerSchema);
+ }
+
+ static List<Integer> nanCountFieldIds(IcebergWriteSchemaContext context,
Schema writerSchema) {
+ return nanCountFieldIds(context.getMetricsConfig(), writerSchema);
+ }
+
+ private static List<Integer> nanCountFieldIds(MetricsConfig metricsConfig,
Schema writerSchema) {
+ return TypeUtil.indexById(writerSchema.asStruct()).values().stream()
+ .filter(field -> isFloatingType(field.type()))
+ .map(Types.NestedField::fieldId)
+ .filter(fieldId -> MetricsUtil.metricsMode(writerSchema,
metricsConfig, fieldId)
+ != MetricsModes.None.get())
+ .sorted()
+ .collect(Collectors.toList());
+ }
+
+ private static boolean isFloatingType(Type type) {
+ return type.typeId() == Type.TypeID.FLOAT || type.typeId() ==
Type.TypeID.DOUBLE;
+ }
}
diff --git
a/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergConnectorTransactionTest.java
b/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergConnectorTransactionTest.java
index 417eee6f0d7..eb30afcddaf 100644
---
a/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergConnectorTransactionTest.java
+++
b/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergConnectorTransactionTest.java
@@ -96,6 +96,10 @@ public class IcebergConnectorTransactionTest {
private static final Schema PART_SCHEMA = new Schema(
Types.NestedField.required(1, "id", Types.IntegerType.get()),
Types.NestedField.required(2, "region", Types.StringType.get()));
+ // A floating identity partition column, the only shape that can carry a
NaN partition value.
+ private static final Schema FLOAT_PART_SCHEMA = new Schema(
+ Types.NestedField.required(1, "id", Types.IntegerType.get()),
+ Types.NestedField.required(2, "d", Types.DoubleType.get()));
private static InMemoryCatalog freshCatalog() {
InMemoryCatalog catalog = new InMemoryCatalog();
@@ -1224,6 +1228,120 @@ public class IcebergConnectorTransactionTest {
"a non-identity (bucket) partition spec disables narrowing ->
the concurrent append still conflicts");
}
+ @Test
+ public void deleteOnNaNIdentityPartitionCommits() {
+ // A NaN cannot be an iceberg literal (Literals.from throws "Cannot
create expression literal from NaN"),
+ // yet parsePartitionValueFromString deliberately turns BE's "nan"
into Double.NaN, so narrowing conflict
+ // detection to the touched partition used to abort the whole DELETE
before validation even ran. iceberg
+ // expresses it as isNaN instead.
+ PartitionSpec spec =
PartitionSpec.builderFor(FLOAT_PART_SCHEMA).identity("d").build();
+ InMemoryCatalog catalog = freshCatalog();
+ TableIdentifier id = TableIdentifier.of("db1", "t1");
+ Table table = catalog.createTable(id, FLOAT_PART_SCHEMA, spec,
props("format-version", "2"));
+ String seedPath = "s3://b/db1/t1/d=NaN/seed.parquet";
+ table.newAppend().appendFile(partitionedDataFile(spec, seedPath, 5L,
"d=NaN")).commit();
+
+ IcebergConnectorTransaction txn =
txnFor(opsReturning(catalog.loadTable(id)), new RecordingConnectorContext());
+ txn.beginWrite(SESSION, "db1", "t1", deleteCtx());
+ TIcebergCommitData del =
positionDeleteItem("s3://b/db1/t1/d=NaN/del.parquet", 2L, seedPath);
+ del.setPartitionValues(Collections.singletonList("nan"));
+ del.setPartitionSpecId(spec.specId());
+ txn.addCommitData(commitBytes(del));
+ txn.commit();
+
+ Snapshot snap = reloadCurrentSnapshot(catalog, id);
+ Assertions.assertEquals("1", snap.summary().get("added-delete-files"),
+ "a delete touching a NaN identity partition must commit, not
fail on the partition literal");
+ }
+
+ @Test
+ public void deleteOnNaNIdentityPartitionStillDetectsConcurrentConflict() {
+ // Not throwing is only half of it: isNaN has to work as a FILTER, so
that narrowing to the NaN
+ // partition still catches a real conflict there. This also pins that
the row-level write constraint
+ // survives the NaN partition predicate -- the whole filter is ANDed
in applyConflictDetectionFilter,
+ // and a dropped NaN arm would silently stop guarding the NaN
partition.
+ PartitionSpec spec =
PartitionSpec.builderFor(FLOAT_PART_SCHEMA).identity("d").build();
+ InMemoryCatalog catalog = freshCatalog();
+ TableIdentifier id = TableIdentifier.of("db1", "t1");
+ Table table = catalog.createTable(id, FLOAT_PART_SCHEMA, spec,
props("format-version", "2"));
+ String seedPath = "s3://b/db1/t1/d=NaN/seed.parquet";
+ table.newAppend().appendFile(partitionedDataFile(spec, seedPath, 5L,
"d=NaN")).commit();
+
+ IcebergConnectorTransaction txn =
txnFor(opsReturning(catalog.loadTable(id)), new RecordingConnectorContext());
+ txn.beginWrite(SESSION, "db1", "t1", deleteCtx());
+ // DELETE ... WHERE d > 5 -- in Doris that is true for a NaN row,
which is why the NaN partition is
+ // the one being deleted from.
+ txn.applyWriteConstraint(new ConnectorPredicate(new
ConnectorComparison(
+ ConnectorComparison.Operator.GT,
+ new ConnectorColumnRef("d", ConnectorType.of("UNKNOWN")),
+ new ConnectorLiteral(ConnectorType.of("DOUBLE"), 5.0d))));
+ // A concurrent append into the very partition being deleted from must
be a conflict.
+ catalog.loadTable(id).newAppend().appendFile(partitionedDataFile(spec,
+ "s3://b/db1/t1/d=NaN/concurrent.parquet", 7L,
"d=NaN")).commit();
+
+ TIcebergCommitData del =
positionDeleteItem("s3://b/db1/t1/d=NaN/del.parquet", 2L, seedPath);
+ del.setPartitionValues(Collections.singletonList("nan"));
+ del.setPartitionSpecId(spec.specId());
+ txn.addCommitData(commitBytes(del));
+
+ DorisConnectorException ex =
Assertions.assertThrows(DorisConnectorException.class, txn::commit,
+ "a concurrent append into the deleted NaN partition must be
detected as a conflict");
+ Assertions.assertTrue(ex.getMessage().contains("is_nan"),
+ "the NaN partition predicate must reach conflict validation,
not be dropped: "
+ + ex.getMessage());
+ }
+
+ @Test
+ public void deleteOnNaNIdentityPartitionExcludesNonMatchingPartition() {
+ // The other side of the same narrowing: a concurrent append into
d=1.0 does not match `d > 5` under
+ // Doris semantics either, so it must NOT be reported as a conflict.
Together with the test above this
+ // proves the NaN predicate narrows rather than degrading to
always-true or always-false.
+ PartitionSpec spec =
PartitionSpec.builderFor(FLOAT_PART_SCHEMA).identity("d").build();
+ InMemoryCatalog catalog = freshCatalog();
+ TableIdentifier id = TableIdentifier.of("db1", "t1");
+ Table table = catalog.createTable(id, FLOAT_PART_SCHEMA, spec,
props("format-version", "2"));
+ String seedPath = "s3://b/db1/t1/d=NaN/seed.parquet";
+ table.newAppend().appendFile(partitionedDataFile(spec, seedPath, 5L,
"d=NaN")).commit();
+
+ IcebergConnectorTransaction txn =
txnFor(opsReturning(catalog.loadTable(id)), new RecordingConnectorContext());
+ txn.beginWrite(SESSION, "db1", "t1", deleteCtx());
+ txn.applyWriteConstraint(new ConnectorPredicate(new
ConnectorComparison(
+ ConnectorComparison.Operator.GT,
+ new ConnectorColumnRef("d", ConnectorType.of("UNKNOWN")),
+ new ConnectorLiteral(ConnectorType.of("DOUBLE"), 5.0d))));
+ catalog.loadTable(id).newAppend().appendFile(partitionedDataFile(spec,
+ "s3://b/db1/t1/d=1.0/concurrent.parquet", 7L,
"d=1.0")).commit();
+
+ TIcebergCommitData del =
positionDeleteItem("s3://b/db1/t1/d=NaN/del.parquet", 2L, seedPath);
+ del.setPartitionValues(Collections.singletonList("nan"));
+ del.setPartitionSpecId(spec.specId());
+ txn.addCommitData(commitBytes(del));
+ txn.commit();
+
+ Snapshot snap = reloadCurrentSnapshot(catalog, id);
+ Assertions.assertEquals("1", snap.summary().get("added-delete-files"),
+ "an append into a partition the DELETE does not match must not
block the commit");
+ }
+
+ @Test
+ public void staticOverwriteIntoNaNPartitionCommits() {
+ // The same literal, reached from INSERT OVERWRITE ...
PARTITION(d='nan') -> buildPartitionFilter.
+ PartitionSpec spec =
PartitionSpec.builderFor(FLOAT_PART_SCHEMA).identity("d").build();
+ InMemoryCatalog catalog = freshCatalog();
+ TableIdentifier id = TableIdentifier.of("db1", "t1");
+ Table table = catalog.createTable(id, FLOAT_PART_SCHEMA, spec,
props("write.format.default", "parquet"));
+ IcebergConnectorTransaction txn = txnFor(opsReturning(table), new
RecordingConnectorContext());
+
+ txn.beginWrite(SESSION, "db1", "t1", overwriteStaticCtx(table,
Collections.singletonMap("d", "nan")));
+
txn.addCommitData(commitBytes(dataFileItem("s3://b/db1/t1/d=NaN/f1.parquet",
4L, 1024L,
+ Collections.singletonList("nan"))));
+ txn.commit();
+
+ Snapshot snap = reloadCurrentSnapshot(catalog, id);
+ Assertions.assertEquals("overwrite", snap.operation());
+ Assertions.assertEquals("1", snap.summary().get("added-data-files"));
+ }
+
@Test
public void isSerializableIsolationLevelDefaultsAndReadsProperty() {
InMemoryCatalog catalog = freshCatalog();
diff --git
a/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergPredicateConverterFloatSemanticsTest.java
b/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergPredicateConverterFloatSemanticsTest.java
new file mode 100644
index 00000000000..b297a37a5ab
--- /dev/null
+++
b/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergPredicateConverterFloatSemanticsTest.java
@@ -0,0 +1,456 @@
+// 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.
+
+package org.apache.doris.connector.iceberg;
+
+import org.apache.doris.connector.spi.ConnectorType;
+import org.apache.doris.connector.spi.pushdown.ConnectorColumnRef;
+import org.apache.doris.connector.spi.pushdown.ConnectorComparison;
+import org.apache.doris.connector.spi.pushdown.ConnectorExpression;
+import org.apache.doris.connector.spi.pushdown.ConnectorIn;
+import org.apache.doris.connector.spi.pushdown.ConnectorLiteral;
+import org.apache.doris.connector.spi.pushdown.ConnectorNot;
+
+import org.apache.iceberg.DataFile;
+import org.apache.iceberg.DataFiles;
+import org.apache.iceberg.FileFormat;
+import org.apache.iceberg.Metrics;
+import org.apache.iceberg.PartitionSpec;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.StructLike;
+import org.apache.iceberg.expressions.And;
+import org.apache.iceberg.expressions.Evaluator;
+import org.apache.iceberg.expressions.Expression;
+import org.apache.iceberg.expressions.InclusiveMetricsEvaluator;
+import org.apache.iceberg.expressions.Or;
+import org.apache.iceberg.types.Conversions;
+import org.apache.iceberg.types.Types;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.nio.ByteBuffer;
+import java.time.ZoneOffset;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * DORIS-29047 (NaN) and its signed-zero sibling: a FLOAT/DOUBLE predicate
pushed to iceberg must not prune a
+ * file that holds rows Doris considers matching. A pruned file never becomes
a split, so BE never sees those
+ * rows and cannot filter them back in — the query silently returns too few
rows rather than failing.
+ *
+ * <p>Two oracles, because either one alone misses the bug:
+ * <ul>
+ * <li>{@link InclusiveMetricsEvaluator} over hand-built file metrics — the
actual pruning decision, and a
+ * direct reproduction of the reported queries (metrics shaped exactly
as spark / iceberg-java write
+ * them: NaN never reaches the bounds, {@code -0.0} stays {@code
-0.0}).</li>
+ * <li>The row-level {@link Evaluator} over the same converted expression —
pins that the expression is
+ * EQUIVALENT to the Doris predicate rather than merely wider.
Equivalence is what makes NOT / AND / OR
+ * compose: iceberg's RewriteNot turns {@code not(or(gt, isNaN))} into
{@code and(ltEq, notNaN)}, which
+ * is exactly Doris's {@code NOT(d > v)} over a NaN row.</li>
+ * </ul>
+ */
+public class IcebergPredicateConverterFloatSemanticsTest {
+
+ private static final int D_ID = 1;
+
+ private static final Schema SCHEMA = new Schema(
+ Types.NestedField.optional(D_ID, "d", Types.DoubleType.get()),
+ Types.NestedField.optional(2, "f", Types.FloatType.get()),
+ Types.NestedField.optional(3, "i", Types.IntegerType.get()));
+
+ // The iceberg spec forbids NaN as a bound ("NaNs are not permitted as
lower or upper bounds"), so a NaN
+ // shows up only in nan_value_counts; -0.0 survives in the bounds because
they use the IEEE total order.
+ private static final DataFile NAN_ONLY = file("nan_only", 1, 1L, null,
null);
+ private static final DataFile NAN_MIXED = file("nan_mixed", 2, 1L, 1.0d,
1.0d);
+ private static final DataFile NEG_ZERO = file("negzero", 1, 0L, -0.0d,
-0.0d);
+ private static final DataFile POS_ZERO = file("poszero", 1, 0L, 0.0d,
0.0d);
+ private static final DataFile PLAIN = file("plain", 2, 0L, 10.0d, 20.0d);
+ // Doris's own writer reports no NaN count at all (IcebergWriterHelper
passes a null nanValueCounts), so
+ // nothing in the metadata can rule NaN out and such a file must survive
every float range predicate.
+ private static final DataFile UNKNOWN_NAN = file("doris_written", 2, null,
1.0d, 1.0d);
+ // Same bounds as NAN_MIXED and UNKNOWN_NAN, but the writer states there
is no NaN -- the only difference
+ // that may bring pruning back.
+ private static final DataFile NO_NAN_ONE_VALUE = file("one_value", 2, 0L,
1.0d, 1.0d);
+ // Every value is NULL: no bounds, no NaN. Doris matches no row of it for
any comparison.
+ private static final DataFile NULL_ONLY = file("null_only", 2, 2L, 0L,
null, null);
+
+ /**
+ * The reported query: {@code WHERE d > 0} / {@code d >= 0} returned
nothing over a single NaN row.
+ *
+ * <p>The four cases here are the deterministic counterpart of the {@code
float_prune_nan_only}
+ * assertions in the {@code test_iceberg_float_predicate_pushdown}
regression suite, whose fixture
+ * carries exactly these metrics.
+ */
+ @Test
+ public void nanOnlyFileSurvivesRangePredicates() {
+ Assertions.assertTrue(mayMatch(single(cmp("d",
ConnectorComparison.Operator.GT, 0.0d)), NAN_ONLY));
+ Assertions.assertTrue(mayMatch(single(cmp("d",
ConnectorComparison.Operator.GE, 0.0d)), NAN_ONLY));
+ // NaN satisfies neither, so the file is still pruned.
+ Assertions.assertFalse(mayMatch(single(cmp("d",
ConnectorComparison.Operator.LT, 0.0d)), NAN_ONLY));
+ Assertions.assertFalse(mayMatch(single(cmp("d",
ConnectorComparison.Operator.EQ, 0.0d)), NAN_ONLY));
+ }
+
+ /**
+ * The half that a "whole file is NaN" special case would miss: NaN hides
outside the bounds, so a file
+ * holding {1.0, NaN} is pruned by the bounds alone for {@code d > 5}.
+ */
+ @Test
+ public void nanMixedFileSurvivesRangePredicateAboveItsBounds() {
+ Assertions.assertTrue(mayMatch(single(cmp("d",
ConnectorComparison.Operator.GT, 5.0d)), NAN_MIXED));
+ Assertions.assertTrue(mayMatch(single(cmp("d",
ConnectorComparison.Operator.GE, 5.0d)), NAN_MIXED));
+ Assertions.assertTrue(mayMatch(single(cmp("d",
ConnectorComparison.Operator.GT, 5.0d)), UNKNOWN_NAN));
+ }
+
+ /**
+ * {@code !=} / {@code NOT IN} prune through {@code uniqueValue}, whose
NaN guard only fires when the file
+ * actually reports a NaN count — so a file written without one is pruned
as if its single bound value
+ * were its only value.
+ */
+ @Test
+ public void notEqualAndNotInSurviveOnFilesThatMayHoldNaN() {
+ Assertions.assertTrue(mayMatch(single(cmp("d",
ConnectorComparison.Operator.NE, 1.0d)), UNKNOWN_NAN));
+ Assertions.assertTrue(mayMatch(single(notIn("d", 1.0d)), UNKNOWN_NAN));
+ }
+
+ /**
+ * The fix must not blanket-disable pruning: a file that reports zero NaNs
is still pruned, which is what
+ * keeps spark / flink / iceberg-java written tables fully prunable.
+ */
+ @Test
+ public void rangePredicatesStillPruneFilesWithoutNaN() {
+ Assertions.assertFalse(mayMatch(single(cmp("d",
ConnectorComparison.Operator.GT, 100.0d)), PLAIN));
+ Assertions.assertFalse(mayMatch(single(cmp("d",
ConnectorComparison.Operator.GE, 100.0d)), PLAIN));
+ Assertions.assertFalse(mayMatch(single(cmp("d",
ConnectorComparison.Operator.LT, 1.0d)), PLAIN));
+ Assertions.assertFalse(mayMatch(single(cmp("d",
ConnectorComparison.Operator.EQ, 1.0d)), PLAIN));
+ Assertions.assertFalse(mayMatch(single(notIn("d", 1.0d)),
NO_NAN_ONE_VALUE));
+ Assertions.assertFalse(
+ mayMatch(single(cmp("d", ConnectorComparison.Operator.NE,
1.0d)), NO_NAN_ONE_VALUE));
+ }
+
+ /**
+ * Signed zero: Doris reads {@code -0.0 == 0.0} (IEEE) while iceberg
orders {@code -0.0} strictly before
+ * {@code +0.0}, so a bound at the wrong zero drops the other one in both
directions.
+ */
+ @Test
+ public void signedZeroFilesSurviveZeroPredicates() {
+ Assertions.assertTrue(mayMatch(single(cmp("d",
ConnectorComparison.Operator.EQ, 0.0d)), NEG_ZERO));
+ Assertions.assertTrue(mayMatch(single(cmp("d",
ConnectorComparison.Operator.GE, 0.0d)), NEG_ZERO));
+ Assertions.assertTrue(mayMatch(single(in("d", 0.0d)), NEG_ZERO));
+ Assertions.assertTrue(mayMatch(single(cmp("d",
ConnectorComparison.Operator.EQ, -0.0d)), POS_ZERO));
+ Assertions.assertTrue(mayMatch(single(cmp("d",
ConnectorComparison.Operator.LE, -0.0d)), POS_ZERO));
+ // Still precise: widening the bound to cover both zeros must not
start keeping unrelated files, and
+ // -0.0 is neither > 0 nor < 0, so both still prune its file (the
regression suite asserts the same
+ // four outcomes over the identically-shaped float_prune_negzero
fixture).
+ Assertions.assertFalse(mayMatch(single(cmp("d",
ConnectorComparison.Operator.GT, 0.0d)), NEG_ZERO));
+ Assertions.assertFalse(mayMatch(single(cmp("d",
ConnectorComparison.Operator.LT, 0.0d)), NEG_ZERO));
+ Assertions.assertFalse(mayMatch(single(cmp("d",
ConnectorComparison.Operator.EQ, 0.0d)), PLAIN));
+ Assertions.assertFalse(mayMatch(single(cmp("d",
ConnectorComparison.Operator.LE, -0.0d)), PLAIN));
+ }
+
+ /**
+ * {@code WHERE d > 0} carries an INT literal, not a double one (Nereids
types the bare {@code 0} as an
+ * integer), so the zero handling has to see through the literal's Java
type.
+ */
+ @Test
+ public void integerZeroLiteralOnFloatingColumnIsStillAZero() {
+ Expression ge = single(new
ConnectorComparison(ConnectorComparison.Operator.GE, col("d"),
+ new ConnectorLiteral(ConnectorType.of("INT"), 0L)));
+ Assertions.assertTrue(mayMatch(ge, NEG_ZERO), "d >= 0 must keep a
-0.0-only file");
+ Assertions.assertTrue(mayMatch(ge, NAN_ONLY), "d >= 0 must keep a
NaN-only file");
+ }
+
+ /** A FLOAT column takes the same path — the reconciliation keys off the
iceberg column type. */
+ @Test
+ public void floatColumnGetsTheSameTreatment() {
+ Expression gt = single(new
ConnectorComparison(ConnectorComparison.Operator.GT,
+ col("f"), new ConnectorLiteral(ConnectorType.of("FLOAT"),
0.0d)));
+ Assertions.assertEquals(Expression.Operation.OR, gt.op());
+ Assertions.assertEquals(Expression.Operation.IS_NAN, ((Or)
gt).right().op());
+ }
+
+ /**
+ * A NaN literal is reachable ({@code WHERE d > cast('nan' as double)} —
Nereids' DoubleLiteral parses the
+ * {@code nan} spelling) and iceberg refuses it outright: {@code
Expressions.*(col, NaN)} throws "Cannot
+ * create expression literal from NaN". It must become the unary
isNaN/notNaN instead of escaping as a
+ * planning failure.
+ */
+ @Test
+ public void nanLiteralMapsToUnaryPredicatesInsteadOfThrowing() {
+ Assertions.assertEquals(Expression.Operation.IS_NAN,
+ single(cmp("d", ConnectorComparison.Operator.EQ,
Double.NaN)).op());
+ Assertions.assertEquals(Expression.Operation.IS_NAN,
+ single(cmp("d", ConnectorComparison.Operator.GE,
Double.NaN)).op());
+ // notNaN, guarded against null (iceberg reads notNaN(null) as true;
Doris does not match it).
+ assertNotNaNAndNotNull(single(cmp("d",
ConnectorComparison.Operator.NE, Double.NaN)));
+ assertNotNaNAndNotNull(single(cmp("d",
ConnectorComparison.Operator.LT, Double.NaN)));
+ Assertions.assertEquals(Expression.Operation.FALSE,
+ single(cmp("d", ConnectorComparison.Operator.GT,
Double.NaN)).op());
+ Assertions.assertEquals(Expression.Operation.NOT_NULL,
+ single(cmp("d", ConnectorComparison.Operator.LE,
Double.NaN)).op());
+ // `d = NaN` keeps exactly the files that may hold a NaN.
+ Assertions.assertTrue(mayMatch(single(cmp("d",
ConnectorComparison.Operator.EQ, Double.NaN)), NAN_ONLY));
+ Assertions.assertFalse(mayMatch(single(cmp("d",
ConnectorComparison.Operator.EQ, Double.NaN)), PLAIN));
+ }
+
+ /**
+ * A NaN literal is the one case where the emitted arm is a bare {@code
notNaN}, and iceberg evaluates
+ * {@code notNaN(null)} as TRUE while Doris leaves {@code NULL != NaN} /
{@code NULL < NaN} UNKNOWN, so a
+ * null row matches nothing. REWRITE turns this expression into the file
filter with no downstream
+ * re-filter, so without the NULL guard a rewrite would pull in a
null-only file holding no matching row.
+ *
+ * <p>Asserted at the file level rather than the row level on purpose:
iceberg's row {@link Evaluator}
+ * compares with a bare {@code Comparator.naturalOrder()} and throws on a
null value, so a null row is not
+ * even expressible there — the file evaluator is what REWRITE actually
consults.</p>
+ */
+ @Test
+ public void nanLiteralNegativeFormsExcludeNullOnlyFiles() {
+ Assertions.assertFalse(mayMatch(single(cmp("d",
ConnectorComparison.Operator.NE, Double.NaN)), NULL_ONLY));
+ Assertions.assertFalse(mayMatch(single(cmp("d",
ConnectorComparison.Operator.LT, Double.NaN)), NULL_ONLY));
+ Assertions.assertFalse(mayMatch(single(notIn("d", Double.NaN)),
NULL_ONLY));
+ Assertions.assertFalse(mayMatch(single(notIn("d", 1.0d, Double.NaN)),
NULL_ONLY));
+ // The positive forms were already null-safe: isNaN and notNull are
both false for a null value.
+ Assertions.assertFalse(mayMatch(single(cmp("d",
ConnectorComparison.Operator.EQ, Double.NaN)), NULL_ONLY));
+ Assertions.assertFalse(mayMatch(single(cmp("d",
ConnectorComparison.Operator.LE, Double.NaN)), NULL_ONLY));
+ // And the guard must not cost anything on a file that does hold
non-null values.
+ Assertions.assertTrue(mayMatch(single(cmp("d",
ConnectorComparison.Operator.NE, Double.NaN)), PLAIN));
+ Assertions.assertTrue(mayMatch(single(cmp("d",
ConnectorComparison.Operator.LT, Double.NaN)), PLAIN));
+ }
+
+ /** IN / NOT IN carry both problems: a listed NaN cannot be an iceberg
literal, a listed zero is two points. */
+ @Test
+ public void inListHandlesNaNAndSignedZero() {
+ Assertions.assertEquals(Expression.Operation.IS_NAN, single(in("d",
Double.NaN)).op());
+ assertNotNaNAndNotNull(single(notIn("d", Double.NaN)));
+ Expression mixed = single(in("d", 1.0d, Double.NaN));
+ Assertions.assertTrue(mayMatch(mixed, NAN_MIXED));
+ Assertions.assertFalse(mayMatch(mixed, PLAIN));
+ }
+
+ /** Non-floating columns keep the plain 1:1 mapping — no isNaN arm, no
zero expansion. */
+ @Test
+ public void nonFloatingColumnsAreUnchanged() {
+ Assertions.assertEquals(Expression.Operation.GT,
+ single(intCmp(ConnectorComparison.Operator.GT)).op());
+ Assertions.assertEquals(Expression.Operation.EQ,
+ single(intCmp(ConnectorComparison.Operator.EQ)).op());
+ Assertions.assertEquals(Expression.Operation.LT_EQ,
+ single(intCmp(ConnectorComparison.Operator.LE)).op());
+ }
+
+ /**
+ * A NOT must not smuggle the bug back in through the FILE evaluator.
Row-level equivalence is not enough
+ * here: iceberg's RewriteNot lowers {@code not(lessThan(d, 5))} to a bare
{@code gtEq(d, 5)}, whose
+ * metrics evaluator prunes a {1.0, NaN} file even though Doris's {@code
NOT(d < 5)} matches the NaN row.
+ * The leaves therefore carry the NaN half on both sides — {@code OR
isNaN} where NaN matches,
+ * {@code AND notNaN} where it does not — so De Morgan turns one into the
other.
+ */
+ @Test
+ public void negatedPredicatesStillKeepFilesHoldingNaN() {
+ // NOT(d < 5) / NOT(d <= 5): Doris matches the NaN row, so the file
must survive.
+ Assertions.assertTrue(mayMatch(single(not(cmp("d",
ConnectorComparison.Operator.LT, 5.0d))), NAN_MIXED));
+ Assertions.assertTrue(mayMatch(single(not(cmp("d",
ConnectorComparison.Operator.LE, 5.0d))), NAN_MIXED));
+ Assertions.assertTrue(mayMatch(single(not(cmp("d",
ConnectorComparison.Operator.LE, 0.0d))), NAN_ONLY));
+ // NOT(d = 1.0) / NOT(d IN (1.0)) over a file whose metrics cannot
rule NaN out.
+ Assertions.assertTrue(mayMatch(single(not(cmp("d",
ConnectorComparison.Operator.EQ, 1.0d))), UNKNOWN_NAN));
+ Assertions.assertTrue(mayMatch(single(not(in("d", 1.0d))),
UNKNOWN_NAN));
+ // Still precise in the other direction: NOT(d = 0) over a -0.0-only
file matches nothing, so the
+ // file is pruned rather than blanket-kept.
+ Assertions.assertFalse(mayMatch(single(not(cmp("d",
ConnectorComparison.Operator.EQ, 0.0d))), NEG_ZERO));
+ // And NOT(d > 100) must still prune a file that lies entirely above
100 with no NaN.
+ Assertions.assertFalse(mayMatch(single(not(cmp("d",
ConnectorComparison.Operator.GE, 1.0d))), PLAIN));
+ }
+
+ /**
+ * The equivalence oracle. iceberg's row-level {@link Evaluator} orders
floats the way Doris does (NaN
+ * greatest), and the signed-zero expansion makes {@code = 0} match both
zeros there too, so the converted
+ * expression must agree with Doris row for row — and its negation must be
the exact complement.
+ */
+ @Test
+ public void convertedExpressionMatchesDorisRowSemantics() {
+ double[] rows = {Double.NaN, -0.0d, 0.0d, 1.0d, -1.0d};
+ double[] lits = {0.0d, -0.0d, 1.0d};
+ List<ConnectorComparison.Operator> ops = Arrays.asList(
+ ConnectorComparison.Operator.EQ,
ConnectorComparison.Operator.NE,
+ ConnectorComparison.Operator.GT,
ConnectorComparison.Operator.GE,
+ ConnectorComparison.Operator.LT,
ConnectorComparison.Operator.LE);
+ for (ConnectorComparison.Operator op : ops) {
+ for (double lit : lits) {
+ Expression expr = single(cmp("d", op, lit));
+ Expression negated = single(not(cmp("d", op, lit)));
+ for (double row : rows) {
+ String where = "d " + op + " " + lit + " over row " + row;
+ boolean expected = dorisMatches(op, row, lit);
+ Assertions.assertEquals(expected, evalRow(expr, row),
where);
+ Assertions.assertEquals(!expected, evalRow(negated, row),
"NOT " + where);
+ }
+ }
+ }
+ }
+
+ // ───────────────────────────────── Doris float semantics (the oracle)
─────────────────────────────────
+
+ // be/src/common/compare.h: NaN equals NaN and is greater than everything
else; zeros compare by IEEE, so
+ // -0.0 == 0.0.
+ private static int dorisCompare(double left, double right) {
+ if (Double.isNaN(left)) {
+ return Double.isNaN(right) ? 0 : 1;
+ }
+ if (Double.isNaN(right)) {
+ return -1;
+ }
+ return Double.compare(left == 0.0d ? 0.0d : left, right == 0.0d ? 0.0d
: right);
+ }
+
+ private static boolean dorisMatches(ConnectorComparison.Operator op,
double row, double lit) {
+ int cmp = dorisCompare(row, lit);
+ switch (op) {
+ case EQ:
+ return cmp == 0;
+ case NE:
+ return cmp != 0;
+ case GT:
+ return cmp > 0;
+ case GE:
+ return cmp >= 0;
+ case LT:
+ return cmp < 0;
+ case LE:
+ return cmp <= 0;
+ default:
+ throw new IllegalArgumentException("unexpected operator " +
op);
+ }
+ }
+
+ // ───────────────────────────────────────────── plumbing
──────────────────────────────────────────────
+
+ private static Expression single(ConnectorExpression expr) {
+ List<Expression> out = new IcebergPredicateConverter(SCHEMA,
ZoneOffset.UTC).convert(expr);
+ Assertions.assertEquals(1, out.size(), "expected exactly one pushed
predicate for " + expr);
+ return out.get(0);
+ }
+
+ private static ConnectorColumnRef col(String name) {
+ return new ConnectorColumnRef(name, ConnectorType.of("UNKNOWN"));
+ }
+
+ private static ConnectorComparison cmp(String colName,
ConnectorComparison.Operator op, double value) {
+ return new ConnectorComparison(op, col(colName),
ConnectorLiteral.ofDouble(value));
+ }
+
+ private static ConnectorNot not(ConnectorExpression operand) {
+ return new ConnectorNot(operand);
+ }
+
+ // `d != NaN` / `d < NaN` / `d NOT IN (NaN)` all lower to notNaN AND
notNull -- the guard exists because
+ // iceberg reads notNaN(null) as TRUE while Doris leaves the comparison
UNKNOWN.
+ private static void assertNotNaNAndNotNull(Expression expr) {
+ Assertions.assertEquals(Expression.Operation.AND, expr.op(),
expr.toString());
+ And and = (And) expr;
+ Assertions.assertEquals(Expression.Operation.NOT_NAN, and.left().op(),
expr.toString());
+ Assertions.assertEquals(Expression.Operation.NOT_NULL,
and.right().op(), expr.toString());
+ }
+
+ private static ConnectorComparison intCmp(ConnectorComparison.Operator op)
{
+ return new ConnectorComparison(op, col("i"), new
ConnectorLiteral(ConnectorType.of("INT"), 0L));
+ }
+
+ private static ConnectorIn in(String colName, double... values) {
+ return connectorIn(colName, false, values);
+ }
+
+ private static ConnectorIn notIn(String colName, double... values) {
+ return connectorIn(colName, true, values);
+ }
+
+ private static ConnectorIn connectorIn(String colName, boolean negated,
double... values) {
+ List<ConnectorExpression> items = new ArrayList<>();
+ for (double value : values) {
+ items.add(ConnectorLiteral.ofDouble(value));
+ }
+ return new ConnectorIn(col(colName), items, negated);
+ }
+
+ private static boolean mayMatch(Expression expr, DataFile dataFile) {
+ return new InclusiveMetricsEvaluator(SCHEMA, expr).eval(dataFile);
+ }
+
+ private static boolean evalRow(Expression expr, double value) {
+ return new Evaluator(SCHEMA.asStruct(), expr).eval(new
DoubleRow(value));
+ }
+
+ private static DataFile file(String name, long records, Long nanCount,
Double lower, Double upper) {
+ return file(name, records, 0L, nanCount, lower, upper);
+ }
+
+ private static DataFile file(String name, long records, long nullCount,
Long nanCount,
+ Double lower, Double upper) {
+ Map<Integer, Long> valueCounts = new HashMap<>();
+ valueCounts.put(D_ID, records);
+ Map<Integer, Long> nullCounts = new HashMap<>();
+ nullCounts.put(D_ID, nullCount);
+ Map<Integer, Long> nanCounts = null;
+ if (nanCount != null) {
+ nanCounts = new HashMap<>();
+ nanCounts.put(D_ID, nanCount);
+ }
+ Metrics metrics = new Metrics(records, null, valueCounts, nullCounts,
nanCounts,
+ bound(lower), bound(upper));
+ return DataFiles.builder(PartitionSpec.unpartitioned())
+ .withPath("/fake/" + name + ".parquet")
+ .withFormat(FileFormat.PARQUET)
+ .withFileSizeInBytes(1024L)
+ .withRecordCount(records)
+ .withMetrics(metrics)
+ .build();
+ }
+
+ private static Map<Integer, ByteBuffer> bound(Double value) {
+ if (value == null) {
+ return null;
+ }
+ Map<Integer, ByteBuffer> bounds = new HashMap<>();
+ bounds.put(D_ID, Conversions.toByteBuffer(Types.DoubleType.get(),
value));
+ return bounds;
+ }
+
+ // Minimal StructLike over SCHEMA; only "d" (position 0) is ever read.
+ private static final class DoubleRow implements StructLike {
+ private final Double value;
+
+ private DoubleRow(Double value) {
+ this.value = value;
+ }
+
+ @Override
+ public int size() {
+ return 3;
+ }
+
+ @Override
+ public <T> T get(int pos, Class<T> javaClass) {
+ return javaClass.cast(pos == 0 ? value : null);
+ }
+
+ @Override
+ public <T> void set(int pos, T newValue) {
+ throw new UnsupportedOperationException();
+ }
+ }
+}
diff --git
a/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergWritePlanProviderTest.java
b/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergWritePlanProviderTest.java
index fd2e4b27d78..17af0421186 100644
---
a/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergWritePlanProviderTest.java
+++
b/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergWritePlanProviderTest.java
@@ -1771,6 +1771,30 @@ public class IcebergWritePlanProviderTest {
return plan.getDataSink().getIcebergMergeSink();
}
+ // The replacement data files an UPDATE / SQL MERGE writes go through the
same iceberg parquet writer as
+ // an INSERT, and VIcebergMergeSink builds its inner TIcebergTableSink by
copying fields across one by
+ // one. A NaN-count policy shipped only on the INSERT sink therefore
leaves merge-written files reporting
+ // no nan_value_counts, so they stay unprunable even when NaN-free — the
pruning restoration would cover
+ // only part of what Doris writes. This pins that the merge dialect
carries it too.
+ @Test
+ public void planWriteMergeSinkShipsTheNanCountFieldIds() {
+ Schema floatSchema = new Schema(
+ Types.NestedField.required(1, "id", Types.IntegerType.get()),
+ Types.NestedField.optional(2, "d", Types.DoubleType.get()));
+ Map<String, String> tableProps = new HashMap<>();
+ tableProps.put("write.format.default", "parquet");
+ tableProps.put("write.data.path", "oss://bucket/wh/db1/t1/data");
+ Table table = freshCatalog().createTable(TableIdentifier.of("db1",
"t1"), floatSchema,
+ PartitionSpec.unpartitioned(), tableProps);
+ TIcebergMergeSink sink = planMergeSink(table, contextWithStorage(),
+ new WriteHandle(new IcebergTableHandle("db1", "t1"))
+ .writeOperation(WriteOperation.MERGE));
+
+ Assertions.assertTrue(sink.isSetNanCountFieldIds(),
+ "the merge dialect must carry the NaN-count policy, not just
the INSERT sink");
+ Assertions.assertEquals(Collections.singletonList(2),
sink.getNanCountFieldIds());
+ }
+
// #66112: UPDATE and SQL MERGE share the TIcebergMergeSink dialect, but
only SQL MERGE carries the
// one-source-row cardinality rule. The engine decides (statement kind)
and the connector must ship the
// decision verbatim: BE gates its duplicate-match validation on this
field, so a connector that dropped
diff --git
a/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergWriterHelperTest.java
b/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergWriterHelperTest.java
index 5d9c3c52218..441ca3d5bc9 100644
---
a/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergWriterHelperTest.java
+++
b/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergWriterHelperTest.java
@@ -238,6 +238,106 @@ public class IcebergWriterHelperTest {
Assertions.assertEquals(Long.valueOf(2L), df.nullValueCounts().get(1));
}
+ /**
+ * NaN counts are the only metadata that can prove a float column holds no
NaN, because iceberg keeps NaN
+ * out of the bounds by spec. A reported zero is therefore a positive
claim that lets a float range
+ * predicate prune the file, so it must travel from BE untouched — and an
absent count must stay absent.
+ */
+ @Test
+ public void convertToWriterResultCarriesNanValueCounts() {
+ Schema floatSchema = new Schema(
+ Types.NestedField.required(1, "id", Types.IntegerType.get()),
+ Types.NestedField.optional(2, "d", Types.DoubleType.get()),
+ Types.NestedField.optional(3, "f", Types.FloatType.get()));
+ Table table = tableWith(floatSchema, "write.format.default",
"parquet");
+
+ TIcebergColumnStats stats = new TIcebergColumnStats();
+ stats.putToValueCounts(2, 10L);
+ stats.putToValueCounts(3, 10L);
+ stats.putToNanValueCounts(2, 3L);
+ // An explicit zero is the whole point: it is what brings pruning back
for a NaN-free column.
+ stats.putToNanValueCounts(3, 0L);
+
+ DataFile df = writeSingle(table, stats, "s3://b/db1/t/f.parquet");
+
+ Assertions.assertEquals(Long.valueOf(3L), df.nanValueCounts().get(2));
+ Assertions.assertEquals(Long.valueOf(0L), df.nanValueCounts().get(3));
+ }
+
+ /**
+ * A BE that reports no NaN counts (an older one, or a format whose writer
cannot count them) must leave the
+ * metric absent, not zero: iceberg reads absent as "may contain NaN" and
keeps the file, whereas a zero
+ * would license pruning away rows that do match.
+ */
+ @Test
+ public void
convertToWriterResultLeavesNanValueCountsUnsetWhenBeReportsNone() {
+ Table table = tableWith("write.format.default", "parquet");
+ TIcebergColumnStats stats = new TIcebergColumnStats();
+ stats.putToValueCounts(1, 10L);
+ stats.putToNullValueCounts(1, 0L);
+
+ DataFile df = writeSingle(table, stats, "s3://b/db1/t/f.parquet");
+
+ Assertions.assertTrue(df.nanValueCounts() == null ||
df.nanValueCounts().isEmpty(),
+ "an unreported NaN count must not become an empty-but-present
or zero claim");
+ }
+
+ /**
+ * Counting NaNs is an extra pass over the data (the parquet footer
carries no NaN count), so BE must not
+ * pay it for a field FE would then drop. {@code shouldCollectColumnStats}
is table-wide and stays true as
+ * soon as any one field is enabled, so it cannot be that gate — these are
the cases where the two differ.
+ */
+ @Test
+ public void nanCountFieldIdsFollowThePerFieldMetricsPolicy() {
+ Schema floatSchema = new Schema(
+ Types.NestedField.required(1, "id", Types.IntegerType.get()),
+ Types.NestedField.optional(2, "d", Types.DoubleType.get()),
+ Types.NestedField.optional(3, "f", Types.FloatType.get()));
+
+ // Default policy: every float field is eligible.
+ Table dflt = tableWith(floatSchema, "write.format.default", "parquet");
+ Assertions.assertEquals(Arrays.asList(2, 3),
+ IcebergWriterHelper.nanCountFieldIds(dflt, floatSchema));
+
+ // default=none with one NON-float field re-enabled: the table-wide
flag is still true, but no float
+ // field survives the policy, so BE must be told to count nothing.
+ Table mixed = tableWith(floatSchema, "write.format.default", "parquet",
+ "write.metadata.metrics.default", "none",
+ "write.metadata.metrics.column.id", "full");
+
Assertions.assertTrue(IcebergWriterHelper.shouldCollectColumnStats(mixed,
floatSchema),
+ "one enabled field keeps table-wide collection on -- which is
exactly why it cannot gate the scan");
+ Assertions.assertTrue(IcebergWriterHelper.nanCountFieldIds(mixed,
floatSchema).isEmpty());
+
+ // A single float field disabled while the rest stay default.
+ Table oneOff = tableWith(floatSchema, "write.format.default",
"parquet",
+ "write.metadata.metrics.column.d", "none");
+ Assertions.assertEquals(Collections.singletonList(3),
+ IcebergWriterHelper.nanCountFieldIds(oneOff, floatSchema));
+ }
+
+ /**
+ * Iceberg gives the default metrics mode only to the first
+ * {@code write.metadata.metrics.max-inferred-column-defaults} fields and
{@code none} to everything after,
+ * so a wide table disables most of its columns with no property set at
all. That is the case that makes
+ * this list worth threading rather than a nice-to-have.
+ */
+ @Test
+ public void nanCountFieldIdsHonorTheInferredColumnCapOnAWideTable() {
+ List<Types.NestedField> fields = new ArrayList<>();
+ for (int id = 1; id <= 6; id++) {
+ fields.add(Types.NestedField.optional(id, "d" + id,
Types.DoubleType.get()));
+ }
+ Schema wide = new Schema(fields);
+ // Cap the inferred defaults at 4 rather than building a 100-column
fixture; the mechanism is the same.
+ Table table = tableWith(wide, "write.format.default", "parquet",
+ "write.metadata.metrics.max-inferred-column-defaults", "4");
+
+
Assertions.assertTrue(IcebergWriterHelper.shouldCollectColumnStats(table,
wide));
+ Assertions.assertEquals(Arrays.asList(1, 2, 3, 4),
+ IcebergWriterHelper.nanCountFieldIds(table, wide),
+ "fields past the inferred-column cap default to metrics=none
and must not be scanned");
+ }
+
// ──────────── convertToWriterResult: #65782 honor iceberg metrics policy
(ported from fe-core) ────────────
@Test
diff --git a/gensrc/thrift/DataSinks.thrift b/gensrc/thrift/DataSinks.thrift
index 036b7e51555..76fbf952100 100644
--- a/gensrc/thrift/DataSinks.thrift
+++ b/gensrc/thrift/DataSinks.thrift
@@ -493,6 +493,12 @@ struct TIcebergTableSink {
17: optional TIcebergWriteType write_type = TIcebergWriteType.INSERT;
// Unset keeps collection enabled for rolling upgrades with older FEs.
18: optional bool collect_column_stats;
+ // Iceberg field ids of the FLOAT/DOUBLE fields whose NaN count would
survive the table's metrics policy
+ // (effective mode != none). Counting a NaN is an extra pass over the data
-- unlike the other statistics,
+ // which the parquet footer already carries -- so BE must not pay it for a
field FE would then drop.
+ // Unset or empty means count nothing: an older FE does not read
nan_value_counts back, so counting for it
+ // would be pure waste, and a table whose float fields are all
metrics-disabled has nothing to report.
+ 19: optional list<i32> nan_count_field_ids;
}
struct TIcebergRewritableDeleteFileSet {
@@ -547,6 +553,10 @@ struct TIcebergMergeSink {
16: optional bool writes_data_files;
// Whether the complete target schema contains Variant; used only to fence
old-BE writer omission.
17: optional bool has_variant_schema;
+ // Same contract as TIcebergTableSink.nan_count_field_ids, computed
against the MERGE schema. The
+ // replacement data files UPDATE / SQL MERGE write go through the same
iceberg parquet writer, so
+ // without this they would report no NaN counts and stay unprunable even
when NaN-free.
+ 18: optional list<i32> nan_count_field_ids;
// delete side (position delete only)
20: optional TFileContent delete_type
diff --git
a/regression-test/data/external_table_p0/iceberg/write/test_iceberg_write_stats2.out
b/regression-test/data/external_table_p0/iceberg/write/test_iceberg_write_stats2.out
index 56233ea9861..de452a61eaa 100644
---
a/regression-test/data/external_table_p0/iceberg/write/test_iceberg_write_stats2.out
+++
b/regression-test/data/external_table_p0/iceberg/write/test_iceberg_write_stats2.out
@@ -7,7 +7,7 @@ true 11 111 1.1 1.1 1111 1234.5678
1234.567890 123456789012345678.123456789012 a
0 PARQUET 2 {1:2, 2:2, 3:2, 4:2, 5:2, 6:2, 7:2, 8:2, 9:2, 10:2,
11:2, 12:2} {1:0, 2:0, 3:0, 4:0, 5:0, 6:0, 7:0, 8:0, 9:0, 10:0, 11:0, 12:0}
{1:0x00, 2:0x0B000000, 3:0x6F00000000000000, 4:0xCDCC8C3F,
5:0x9A9999999999F13F, 6:0x0457, 7:0x00BC614E, 8:0x499602D2,
9:0x018EE90FF6C373E0393713FA14, 10:0x616161, 11:0x9E4B0000,
12:0x005CE70F33F10500} {1:0x01, 2:0x16000000, 3:0xDE00000000000000,
4:0xCDCC0C40, 5:0x9A99999999990140, 6:0x08AE, 7:0x05397FB1, 8:0x020A75E124,
9:0x0C7748819DFFB62505316873C [...]
-- !sql_2 --
-{"bigint_col":{"column_size":74, "value_count":2, "null_value_count":0,
"nan_value_count":null, "lower_bound":111, "upper_bound":222},
"boolean_col":{"column_size":33, "value_count":2, "null_value_count":0,
"nan_value_count":null, "lower_bound":0, "upper_bound":1},
"date_col":{"column_size":66, "value_count":2, "null_value_count":0,
"nan_value_count":null, "lower_bound":"2023-01-01",
"upper_bound":"2023-06-15"}, "datetime_col1":{"column_size":74,
"value_count":2, "null_value_count":0, "n [...]
+{"bigint_col":{"column_size":74, "value_count":2, "null_value_count":0,
"nan_value_count":null, "lower_bound":111, "upper_bound":222},
"boolean_col":{"column_size":33, "value_count":2, "null_value_count":0,
"nan_value_count":null, "lower_bound":0, "upper_bound":1},
"date_col":{"column_size":66, "value_count":2, "null_value_count":0,
"nan_value_count":null, "lower_bound":"2023-01-01",
"upper_bound":"2023-06-15"}, "datetime_col1":{"column_size":74,
"value_count":2, "null_value_count":0, "n [...]
-- !sql_3 --
false 22 222 2.2 2.2 2222 8765.4321 8765.432100
987654321098765432.987654321099 bbb 2023-06-15 2023-06-15T23:45:01
@@ -17,7 +17,7 @@ true 11 111 1.1 1.1 1111 1234.5678
1234.567890 123456789012345678.123456789012 a
0 ORC 2 {1:2, 2:2, 3:2, 4:2, 5:2, 6:2, 7:2, 8:2, 9:2, 10:2,
11:2, 12:2} {} {1:0x00, 2:0x0B00000000000000, 3:0x6F00000000000000,
4:0xCDCC8C3F, 5:0x9A9999999999F13F, 6:0x0457, 7:0x00BC614E, 8:0x075BCD15,
9:0x018EE90FF6C373E0393713FA14, 10:0x616161, 11:0x9E4B0000,
12:0x005CE70F33F10500} {1:0x01, 2:0x1600000000000000,
3:0xDE00000000000000, 4:0xCDCC0C40, 5:0x9A99999999990140, 6:0x08AE,
7:0x05397FB1, 8:0x05397FB1, 9:0x0C7748819DFFB62505316873CB, 10:0x626262,
11:0x434C0000, 12:0x40D91FA833FE0500}
-- !sql_5 --
-{"bigint_col":{"column_size":null, "value_count":2, "null_value_count":null,
"nan_value_count":null, "lower_bound":111, "upper_bound":222},
"boolean_col":{"column_size":null, "value_count":2, "null_value_count":null,
"nan_value_count":null, "lower_bound":0, "upper_bound":1},
"date_col":{"column_size":null, "value_count":2, "null_value_count":null,
"nan_value_count":null, "lower_bound":"2023-01-01",
"upper_bound":"2023-06-15"}, "datetime_col1":{"column_size":null,
"value_count":2, "null_v [...]
+{"bigint_col":{"column_size":null, "value_count":2, "null_value_count":null,
"nan_value_count":null, "lower_bound":111, "upper_bound":222},
"boolean_col":{"column_size":null, "value_count":2, "null_value_count":null,
"nan_value_count":null, "lower_bound":0, "upper_bound":1},
"date_col":{"column_size":null, "value_count":2, "null_value_count":null,
"nan_value_count":null, "lower_bound":"2023-01-01",
"upper_bound":"2023-06-15"}, "datetime_col1":{"column_size":null,
"value_count":2, "null_v [...]
-- !sql_6 --
{1:0x00, 2:0x0B000000, 3:0x6F00000000000000, 4:0xCDCC8C3F,
5:0x9A9999999999F13F, 6:0x0457, 7:0x00BC614E, 8:0x499602D2,
9:0x018EE90FF6C373E0393713FA14, 10:0x616161, 11:0x9E4B0000,
12:0x005CE70F33F10500}
diff --git
a/regression-test/suites/external_table_p0/iceberg/test_iceberg_float_predicate_pushdown.groovy
b/regression-test/suites/external_table_p0/iceberg/test_iceberg_float_predicate_pushdown.groovy
new file mode 100644
index 00000000000..27a32ede55f
--- /dev/null
+++
b/regression-test/suites/external_table_p0/iceberg/test_iceberg_float_predicate_pushdown.groovy
@@ -0,0 +1,124 @@
+// 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.
+
+// Float/double predicate pushdown must not prune iceberg files that hold rows
Doris considers matching.
+// Doris matches rows with "NaN is greater than everything" and IEEE zero
equality (-0.0 == 0.0), while
+// iceberg prunes files with Double.compare, where NaN is absent from the
bounds entirely and -0.0 sorts
+// strictly before +0.0. A pruned file never becomes a split, so the rows are
simply missing from the
+// result with no error -- which is why every assertion here is on
inputSplitNum: it is the exact
+// observable the bug moves, and 0 vs 1 is the difference between losing rows
and returning them.
+//
+// Fixtures are created by spark in run32.sql (one data file each, verified
metadata):
+// float_prune_nan_only nan_value_count=1, no bounds (NaN
never reaches the bounds)
+// float_prune_nan_mixed nan_value_count=1, bounds [1.0, 1.0] (NaN
hides outside the bounds)
+// float_prune_nan_mixed_float same, on a FLOAT column
+// float_prune_negzero nan_value_count=0, bounds [-0.0, -0.0]
+// They must be spark-written: Doris reports no nan_value_counts at all and
normalizes zero bounds, so
+// Doris-written files do not carry the metadata shape that triggers the
pruning.
+suite("test_iceberg_float_predicate_pushdown", "p0,external") {
+ String enabled = context.config.otherConfigs.get("enableIcebergTest")
+ if (enabled == null || !enabled.equalsIgnoreCase("true")) {
+ logger.info("disable iceberg test.")
+ return
+ }
+
+ String rest_port = context.config.otherConfigs.get("iceberg_rest_uri_port")
+ String minio_port = context.config.otherConfigs.get("iceberg_minio_port")
+ String externalEnvIp = context.config.otherConfigs.get("externalEnvIp")
+ String catalog_name = "test_iceberg_float_predicate_pushdown"
+
+ sql """drop catalog if exists ${catalog_name}"""
+ sql """CREATE CATALOG ${catalog_name} PROPERTIES (
+ 'type'='iceberg',
+ 'iceberg.catalog.type'='rest',
+ 'uri' = 'http://${externalEnvIp}:${rest_port}',
+ "s3.access_key" = "admin",
+ "s3.secret_key" = "password",
+ "s3.endpoint" = "http://${externalEnvIp}:${minio_port}",
+ "s3.region" = "us-east-1"
+ );"""
+
+ // Batch mode defers split generation, which makes inputSplitNum in
EXPLAIN non-deterministic.
+ sql """set enable_external_table_batch_mode=false"""
+ sql """switch ${catalog_name}"""
+ sql """use test_db"""
+
+ // ---- NaN, a file that is entirely NaN
-------------------------------------------------------------
+ // Before the fix both of these reported inputSplitNum=0 while `select d >
0` projected true.
+ explain {
+ sql("select id from float_prune_nan_only where d > 0")
+ contains "inputSplitNum=1"
+ }
+ explain {
+ sql("select id from float_prune_nan_only where d >= 0")
+ contains "inputSplitNum=1"
+ }
+ // NaN satisfies neither, so the file is still pruned -- pruning must not
be disabled wholesale.
+ explain {
+ sql("select id from float_prune_nan_only where d < 0")
+ contains "inputSplitNum=0"
+ }
+ explain {
+ sql("select id from float_prune_nan_only where d = 0")
+ contains "inputSplitNum=0"
+ }
+
+ // ---- NaN mixed with ordinary values
---------------------------------------------------------------
+ // The half a "whole file is NaN" special case would miss: bounds are
[1.0, 1.0], so `d > 5` prunes the
+ // file on the bounds alone and the NaN row is lost.
+ explain {
+ sql("select id from float_prune_nan_mixed where d > 5")
+ contains "inputSplitNum=1"
+ }
+ explain {
+ sql("select id from float_prune_nan_mixed where d >= 5")
+ contains "inputSplitNum=1"
+ }
+ // NOT over a range: iceberg's RewriteNot lowers NOT(d < 5) back to a bare
`d >= 5`, so this needs the
+ // NaN arm just as much as the direct comparison does.
+ explain {
+ sql("select id from float_prune_nan_mixed where not (d < 5)")
+ contains "inputSplitNum=1"
+ }
+ explain {
+ sql("select id from float_prune_nan_mixed_float where f > 5")
+ contains "inputSplitNum=1"
+ }
+
+ // ---- Signed zero
----------------------------------------------------------------------------------
+ // The file holds only -0.0. Doris reads -0.0 = 0.0, iceberg orders -0.0
before +0.0, so a bound placed
+ // at +0.0 prunes it.
+ explain {
+ sql("select id from float_prune_negzero where d = 0")
+ contains "inputSplitNum=1"
+ }
+ explain {
+ sql("select id from float_prune_negzero where d >= 0")
+ contains "inputSplitNum=1"
+ }
+ // -0.0 is neither > 0 nor < 0, so both still prune.
+ explain {
+ sql("select id from float_prune_negzero where d > 0")
+ contains "inputSplitNum=0"
+ }
+ explain {
+ sql("select id from float_prune_negzero where d < 0")
+ contains "inputSplitNum=0"
+ }
+
+ sql """drop catalog if exists ${catalog_name}"""
+}
diff --git
a/regression-test/suites/external_table_p0/iceberg/test_iceberg_write_nan_value_counts.groovy
b/regression-test/suites/external_table_p0/iceberg/test_iceberg_write_nan_value_counts.groovy
new file mode 100644
index 00000000000..6a091b838fb
--- /dev/null
+++
b/regression-test/suites/external_table_p0/iceberg/test_iceberg_write_nan_value_counts.groovy
@@ -0,0 +1,109 @@
+// 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.
+
+// Iceberg keeps NaN out of a column's bounds by spec, so nan_value_counts is
the only metadata that can
+// prove a float column holds no NaN. Without it, a float range predicate has
to keep every file (see
+// test_iceberg_float_predicate_pushdown for why). BE now counts NaNs while
writing and FE forwards them,
+// which is what lets a Doris-written file be pruned again.
+//
+// The whole chain is asserted through inputSplitNum, end to end: BE counting
-> TIcebergColumnStats ->
+// IcebergWriterHelper metrics -> manifest -> InclusiveMetricsEvaluator.
+suite("test_iceberg_write_nan_value_counts", "p0,external") {
+ String enabled = context.config.otherConfigs.get("enableIcebergTest")
+ if (enabled == null || !enabled.equalsIgnoreCase("true")) {
+ logger.info("disable iceberg test.")
+ return
+ }
+
+ String rest_port = context.config.otherConfigs.get("iceberg_rest_uri_port")
+ String minio_port = context.config.otherConfigs.get("iceberg_minio_port")
+ String externalEnvIp = context.config.otherConfigs.get("externalEnvIp")
+ String catalog_name = "test_iceberg_write_nan_value_counts"
+ String db_name = catalog_name + "_db"
+
+ sql """drop catalog if exists ${catalog_name}"""
+ sql """CREATE CATALOG ${catalog_name} PROPERTIES (
+ 'type'='iceberg',
+ 'iceberg.catalog.type'='rest',
+ 'uri' = 'http://${externalEnvIp}:${rest_port}',
+ "s3.access_key" = "admin",
+ "s3.secret_key" = "password",
+ "s3.endpoint" = "http://${externalEnvIp}:${minio_port}",
+ "s3.region" = "us-east-1"
+ );"""
+
+ sql """drop database if exists ${catalog_name}.${db_name} force"""
+ sql """create database ${catalog_name}.${db_name}"""
+ // Batch mode defers split generation, which makes inputSplitNum in
EXPLAIN non-deterministic.
+ sql """set enable_external_table_batch_mode=false"""
+ sql """use ${catalog_name}.${db_name}"""
+
+ sql """drop table if exists write_nan_finite"""
+ sql """create table write_nan_finite (id int, d double)"""
+ sql """insert into write_nan_finite values (1, 1.0), (2, 2.0)"""
+
+ sql """drop table if exists write_nan_present"""
+ sql """create table write_nan_present (id int, d double)"""
+ sql """insert into write_nan_present values (1, 1.0), (2, cast('nan' as
double))"""
+
+ // A file BE counted and found NaN-free reports zero, which is what brings
pruning back: the bounds are
+ // [1.0, 2.0] and iceberg can now rule the file out. Before BE reported
NaN counts this was
+ // inputSplitNum=1, because an unknown count forced the file to be kept.
+ explain {
+ sql("select id from write_nan_finite where d > 100")
+ contains "inputSplitNum=0"
+ }
+ explain {
+ sql("select id from write_nan_finite where d >= 100")
+ contains "inputSplitNum=0"
+ }
+ // Pruning must stay honest in the other direction: the NaN row satisfies
`d > 100` in Doris, so its
+ // file carries a positive NaN count and must still be read.
+ explain {
+ sql("select id from write_nan_present where d > 100")
+ contains "inputSplitNum=1"
+ }
+
+ // A predicate the bounds themselves satisfy is unaffected by either count.
+ explain {
+ sql("select id from write_nan_finite where d > 0")
+ contains "inputSplitNum=1"
+ }
+
+ // ---- The same, written as ORC
-----------------------------------------------------------------
+ // ORC column statistics carry no NaN count either, and a Doris-written
ORC file does report bounds,
+ // so without counting it the `OR isNaN` arm would keep every NaN-free ORC
file that the bounds alone
+ // used to prune. Both writers now feed the same counter.
+ sql """drop table if exists write_nan_finite_orc"""
+ sql """create table write_nan_finite_orc (id int, d double) properties
("write-format"="orc")"""
+ sql """insert into write_nan_finite_orc values (1, 1.0), (2, 2.0)"""
+
+ sql """drop table if exists write_nan_present_orc"""
+ sql """create table write_nan_present_orc (id int, d double) properties
("write-format"="orc")"""
+ sql """insert into write_nan_present_orc values (1, 1.0), (2, cast('nan'
as double))"""
+
+ explain {
+ sql("select id from write_nan_finite_orc where d > 100")
+ contains "inputSplitNum=0"
+ }
+ explain {
+ sql("select id from write_nan_present_orc where d > 100")
+ contains "inputSplitNum=1"
+ }
+
+ sql """drop catalog if exists ${catalog_name}"""
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]