Skip to content
Open
2 changes: 2 additions & 0 deletions docs/en/engines/table-engines/integrations/iceberg.md
Original file line number Diff line number Diff line change
Expand Up @@ -98,6 +98,8 @@ The following table shows how Iceberg data types are mapped to ClickHouse data t
| `fixed(N)` | `FixedString(N)` | |
| `decimal(P, S)` | `Decimal(P, S)` | |

Writes use the inverse mapping: `DateTime64(9)` is stored as `timestamp_ns`, and `DateTime64` with an explicit timezone and scale 9 as `timestamptz_ns`, but only when `iceberg_format_version` is 3 or higher. `DateTime` and `DateTime64` with scale 6 or less are stored as `timestamp`. Creating or altering an Iceberg table with `DateTime64` scale 9 on format version 1 or 2 raises an exception.

### Complex types {#complex-types}

| Iceberg type | ClickHouse type |
Expand Down
4 changes: 3 additions & 1 deletion src/Storages/ObjectStorage/DataLakes/Iceberg/Compaction.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -987,7 +987,9 @@ static void writeMetadataFiles(
auto log = getLogger("IcebergCompaction");

ColumnsDescription columns_description = ColumnsDescription::fromNamesAndTypes(sample_block_->getNamesAndTypes());
auto [metadata_object, metadata_object_str] = createEmptyMetadataFile(table_path, columns_description, nullptr, nullptr, context);
const auto format_version = plan.initial_metadata_object->getValue<UInt64>(Iceberg::f_format_version);
auto [metadata_object, metadata_object_str] = createEmptyMetadataFile(
table_path, columns_description, nullptr, nullptr, context, format_version);

auto current_schema_id = metadata_object->getValue<Int64>(Iceberg::f_current_schema_id);
Poco::JSON::Object::Ptr current_schema;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
#include <cctype>
#include <limits>

#include <base/unaligned.h>
#include <Columns/IColumn.h>
#include <Common/Exception.h>
#include <DataTypes/DataTypeNullable.h>
Expand All @@ -22,6 +23,68 @@ extern const int BAD_ARGUMENTS;
namespace Iceberg
{

namespace
{

/// 1e16 µs is ~year 2286. That is not an Iceberg `timestamp` max: Spark sentinels
/// such as `9999-12-31` (`~2.53e17` µs) and ClickHouse `DateTime64` max
/// `2299-12-31` (`~1.04e16` µs) are spec-correct microseconds above this.
/// 1e18 ns is ~year 2001; customer 2026 ns bounds are `~1.76e18`.
constexpr Int64 microseconds_ambiguous_threshold = 10'000'000'000'000'000LL;
constexpr Int64 nanoseconds_lower_threshold = 1'000'000'000'000'000'000LL;
constexpr Int64 nanoseconds_per_microsecond = 1000;

bool magnitudeExceeds(Int64 value, Int64 threshold)
{
return value > threshold || value < -threshold;
}

/// Convert ns ticks to us ticks without shrinking a [lower, upper] range.
Int64 nanosecondsToMicrosecondsForBound(Int64 nanoseconds, bool lower_bound)
{
Int64 microseconds = nanoseconds / nanoseconds_per_microsecond;
if (nanoseconds % nanoseconds_per_microsecond == 0)
return microseconds;

/// C++ integer division truncates toward zero.
if (lower_bound)
{
if (nanoseconds < 0)
--microseconds;
}
else if (nanoseconds > 0)
{
++microseconds;
}
return microseconds;
}

/// Iceberg timestamps are 8-byte little-endian (spec Appendix D).
/// Some writers store nanosecond stats on Iceberg `timestamp` (`DateTime64(6)`).
/// Convert only in a true nanosecond band (`|ticks| > 1e18`). Between 1e16 and
/// 1e18 the value may be far-future microseconds (do not convert; fail open so
/// min/max pruning is skipped). `|ticks| <= 1e16` is kept as microseconds.
std::optional<Field> deserializeDateTime64FromBinaryRepr(const std::string & str, const IDataType & type, bool lower_bound)
{
if (str.size() != sizeof(Int64))
return std::nullopt;

Int64 unscaled = unalignedLoadLittleEndian<Int64>(str.data());
const UInt32 scale = getDecimalScale(type);

if (scale == 6 && magnitudeExceeds(unscaled, microseconds_ambiguous_threshold))
{
if (!magnitudeExceeds(unscaled, nanoseconds_lower_threshold))
return std::nullopt;

unscaled = nanosecondsToMicrosecondsForBound(unscaled, lower_bound);
Comment on lines +75 to +80

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Skip pruning instead of guessing nanosecond bounds

When an Iceberg timestamp manifest contains a magnitude above this threshold, the schema still identifies the value as microseconds; magnitude alone cannot prove that the writer actually used nanoseconds. Dividing the value and then trusting the resulting range allows mayBeTrueInRange to prune a file whose corrupt or differently encoded bounds do not cover its data, producing missing query results. Treat these suspicious bounds as unavailable and skip min/max pruning rather than taking a consequential fallback path.

AGENTS.md reference: AGENTS.md:L153-L153

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is not a P1. The overlapping band already fail-opens; converting |ticks| > 1e18 is the customer ns-stats case, not a guess inside the µs/ns overlap.
Iceberg timestamp is still microseconds in the schema. Magnitude cannot prove nanoseconds, which is why |ticks| in (1e16, 1e18] already returns nullopt. Callers then skip that hyperrectangle (if (!left || !right) continue) and min/max pruning for that column is skipped. That covers spec-correct far-future microseconds: Spark 9999-12-31 (~2.53e17 µs) and ClickHouse DateTime64 max 2299-12-31 (~1.04e16 µs).
Conversion runs only for |ticks| > 1e18. As microseconds that is ~year 33658, which is outside ClickHouse DateTime64(6) (max ~2299, ~1.04e16 µs) and outside normal Iceberg timestamp sentinels. As nanoseconds it is the customer layout (~1.76e18 for 2026) that previously over-pruned every file.
Fail-opening that band as well would keep query results correct (files scanned) but would disable min/max pruning again for those tables (IcebergMinMaxIndexPrunedFiles: 0). That undoes the point of this heuristic.
A corrupt 8-byte value > 1e18 that is not nanoseconds could still be converted into a fake in-range window. That is a residual, not a realistic P1 on this path.

}

return DecimalField<DateTime64>(DateTime64(unscaled), scale);
}

}

Int64 fieldToInt64(const Field & value, std::string_view context, std::string_view arg_name)
{
if (value.getType() == Field::Types::Int64)
Expand Down Expand Up @@ -153,6 +216,9 @@ static std::optional<Field> deserializeDecimalBound(const std::string & str, UIn
std::optional<Field> deserializeFieldFromBinaryRepr(const std::string & str, DataTypePtr expected_type, bool lower_bound)
{
auto non_nullable_type = removeNullable(expected_type);
if (WhichDataType(non_nullable_type).isDateTime64())
return deserializeDateTime64FromBinaryRepr(str, *non_nullable_type, lower_bound);

auto column = non_nullable_type->createColumn();
if (WhichDataType(non_nullable_type).isDecimal())
{
Expand Down
22 changes: 21 additions & 1 deletion src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergWrites.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
#include <Core/Range.h>
#include <Core/Settings.h>
#include <Core/TypeId.h>
#include <DataTypes/DataTypeDateTime64.h>
#include <DataTypes/DataTypeNullable.h>
#include <DataTypes/DataTypeString.h>
#include <DataTypes/DataTypeTime64.h>
Expand Down Expand Up @@ -210,6 +211,25 @@ Int64 getTimeValueInMicroseconds(const Field & field, DataTypePtr type)
throw Exception(ErrorCodes::LOGICAL_ERROR, "Expected Time or Time64, got {}", type->getName());
}

Int64 getDateTime64ValueForIcebergBounds(const Field & field, DataTypePtr type)
{
if (type->isNullable())
return getDateTime64ValueForIcebergBounds(field, assert_cast<const DataTypeNullable *>(type.get())->getNestedType());

if (!WhichDataType(type).isDateTime64())
throw Exception(ErrorCodes::LOGICAL_ERROR, "Expected DateTime64, got {}", type->getName());

const auto scale = getDecimalScale(*type);
const auto value = field.safeGet<Decimal64>().getValue().value;
if (scale == 9)
return value;
if (scale > 6)
throw Exception(ErrorCodes::BAD_ARGUMENTS, "Unsupported type for iceberg {}", type->getName());

/// Iceberg `timestamp` bounds are microseconds. DateTime64(s) for s < 6 stores coarser ticks.
return value * DataTypeDateTime64::getScaleMultiplier(6 - scale).value;
}

template <typename DecimalType>
std::vector<uint8_t> dumpDecimalValue(const Field & field)
{
Expand Down Expand Up @@ -321,7 +341,7 @@ std::vector<uint8_t> dumpFieldToBytes(const Field & field, DataTypePtr type)
case TypeIndex::UInt64:
return dumpValue(applyVisitor(FieldVisitorConvertToNumber<Int64>(), field));
case TypeIndex::DateTime64:
return dumpValue(field.safeGet<Decimal64>().getValue().value);
return dumpValue(getDateTime64ValueForIcebergBounds(field, type));
case TypeIndex::String:
{
auto value = field.safeGet<String>();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,37 @@ namespace DB::ErrorCodes
extern const int LOGICAL_ERROR;
}

namespace
{

/// Avro does not map timestamp-nanos to DateTime64, so identity partition values for
/// Iceberg `timestamp_ns` / `timestamptz_ns` (and similarly `Time64`) arrive as Int64.
/// `KeyCondition` compares DateTime64/Time64 predicates as DecimalField with a scale,
/// so wrap the raw ticks before `mayBeTrueInRange`. Calendar transforms stay Int64.
void wrapIntegerPartitionValuesAsDateTime64(Row & partition_values, const DataTypes & data_types)
{
const size_t size = std::min(partition_values.size(), data_types.size());
for (size_t i = 0; i < size; ++i)
{
auto & field = partition_values[i];
if (field.isNull() || field.getType() != Field::Types::Int64)
continue;

auto type = removeNullable(data_types[i]);
WhichDataType which(type);
if (which.isDateTime64())
{
field = DecimalField<DateTime64>(DateTime64(field.safeGet<Int64>()), getDecimalScale(*type));
}
else if (which.isTime64())
{
field = DecimalField<Time64>(Time64(field.safeGet<Int64>()), getDecimalScale(*type));
}
}
}

}

namespace DB::Iceberg
{

Expand Down Expand Up @@ -246,7 +277,9 @@ PruningReturnStatus ManifestFilesPruner::canBePruned(
{
if (partition_key_condition.has_value())
{
const auto & partition_value = entry->parsed_entry->partition_key_value;
Row partition_value = entry->parsed_entry->partition_key_value;
wrapIntegerPartitionValuesAsDateTime64(partition_value, partition_key->data_types);

std::vector<FieldRef> index_value(partition_value.begin(), partition_value.end());
for (size_t i = 0; i < index_value.size(); ++i)
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -307,7 +307,8 @@ bool MetadataGenerator::isAddColumnApplied(const String & column_name, DataTypeP
return false;

Int32 unused_field_id = metadata_object->getValue<Int32>(Iceberg::f_last_column_id);
auto expected_type = Iceberg::getIcebergType(type, unused_field_id);
auto format_version = metadata_object->getValue<UInt64>(Iceberg::f_format_version);
auto expected_type = Iceberg::getIcebergType(type, unused_field_id, format_version);

auto fields = current_schema->getArray(Iceberg::f_fields);
for (UInt32 i = 0; i < fields->size(); ++i)
Expand Down Expand Up @@ -364,7 +365,8 @@ bool MetadataGenerator::isModifyColumnApplied(const String & column_name, DataTy
return false;

Int32 unused_field_id = metadata_object->getValue<Int32>(Iceberg::f_last_column_id);
auto expected_type = Iceberg::getIcebergType(type, unused_field_id);
auto format_version = metadata_object->getValue<UInt64>(Iceberg::f_format_version);
auto expected_type = Iceberg::getIcebergType(type, unused_field_id, format_version);

auto fields = current_schema->getArray(Iceberg::f_fields);
for (UInt32 i = 0; i < fields->size(); ++i)
Expand Down Expand Up @@ -739,7 +741,8 @@ void MetadataGenerator::generateAddColumnMetadata(const String & column_name, Da

auto last_column_id = metadata_object->getValue<Int32>(Iceberg::f_last_column_id);

auto new_type = Iceberg::getIcebergType(type, last_column_id);
auto format_version = metadata_object->getValue<UInt64>(Iceberg::f_format_version);
auto new_type = Iceberg::getIcebergType(type, last_column_id, format_version);
Poco::JSON::Object::Ptr new_field = new Poco::JSON::Object;
new_field->set(Iceberg::f_id, last_column_id + 1);
new_field->set(Iceberg::f_name, column_name);
Expand All @@ -759,7 +762,9 @@ bool MetadataGenerator::generateModifyColumnMetadata(const String & column_name,
auto current_schema = getCurrentSchema();

auto last_column_id = metadata_object->getValue<Int32>(Iceberg::f_last_column_id);
auto new_type = Iceberg::getIcebergType(type, last_column_id);
auto format_version = metadata_object->getValue<UInt64>(Iceberg::f_format_version);

auto new_type = Iceberg::getIcebergType(type, last_column_id, format_version);
auto schema_fields = current_schema->getArray(Iceberg::f_fields);

for (UInt32 i = 0; i < schema_fields->size(); ++i)
Expand Down
46 changes: 36 additions & 10 deletions src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -535,7 +535,7 @@ static size_t icebergDecimalRequiredBytes(UInt32 precision)
}

/// Returns type and required
std::pair<Poco::Dynamic::Var, bool> getIcebergType(DataTypePtr type, Int32 & iter)
std::pair<Poco::Dynamic::Var, bool> getIcebergType(DataTypePtr type, Int32 & iter, UInt64 format_version)
{
switch (type->getTypeId())
{
Expand All @@ -562,8 +562,30 @@ std::pair<Poco::Dynamic::Var, bool> getIcebergType(DataTypePtr type, Int32 & ite
case TypeIndex::Date32:
return {"date", true};
case TypeIndex::DateTime:
case TypeIndex::DateTime64:
return {"timestamp", true};
case TypeIndex::DateTime64:
{
/// Iceberg `timestamp` is microseconds; `timestamp_ns` is nanoseconds and is
/// allowed only in format version 3. Writing it into v1/v2 metadata is invalid.
const auto scale = getDecimalScale(*type);
if (scale == 9)
{
if (format_version < 3)
throw Exception(
ErrorCodes::BAD_ARGUMENTS,
"Iceberg type {} requires format version 3 or higher, got version {}",
type->getName(),
format_version);

auto date_time64 = std::static_pointer_cast<const DataTypeDateTime64>(type);
if (date_time64->hasExplicitTimeZone())
return {Iceberg::f_timestamptz_ns, true};
return {Iceberg::f_timestamp_ns, true};
}
if (scale <= 6)
return {Iceberg::f_timestamp, true};
throw Exception(ErrorCodes::BAD_ARGUMENTS, "Unsupported type for iceberg {}", type->getName());
}
case TypeIndex::Time:
return {"time", true};
case TypeIndex::Time64:
Expand Down Expand Up @@ -602,7 +624,7 @@ std::pair<Poco::Dynamic::Var, bool> getIcebergType(DataTypePtr type, Int32 & ite
Poco::JSON::Object::Ptr field = new Poco::JSON::Object;
field->set(Iceberg::f_id, ++iter_fields);
field->set(Iceberg::f_name, type_tuple->getNameByPosition(iter_names));
auto child_type = getIcebergType(element->getNormalizedType(), iter);
auto child_type = getIcebergType(element->getNormalizedType(), iter, format_version);
field->set(Iceberg::f_required, child_type.second);
field->set(Iceberg::f_type, child_type.first);
fields->add(field);
Expand All @@ -618,7 +640,7 @@ std::pair<Poco::Dynamic::Var, bool> getIcebergType(DataTypePtr type, Int32 & ite

field->set(Iceberg::f_type, "list");
field->set(Iceberg::f_element_id, ++iter);
auto child_type = getIcebergType(type_array->getNestedType(), iter);
auto child_type = getIcebergType(type_array->getNestedType(), iter, format_version);
field->set(Iceberg::f_required, false);
field->set(Iceberg::f_element, child_type.first);
field->set(Iceberg::f_element_required, child_type.second);
Expand All @@ -633,16 +655,16 @@ std::pair<Poco::Dynamic::Var, bool> getIcebergType(DataTypePtr type, Int32 & ite
field->set(Iceberg::f_key_id, ++iter);
field->set(Iceberg::f_value_id, ++iter);

field->set(Iceberg::f_key, getIcebergType(type_map->getKeyType(), iter).first);
auto value_type = getIcebergType(type_map->getValueType(), iter);
field->set(Iceberg::f_key, getIcebergType(type_map->getKeyType(), iter, format_version).first);
auto value_type = getIcebergType(type_map->getValueType(), iter, format_version);
field->set(Iceberg::f_value, value_type.first);
field->set(Iceberg::f_value_required, value_type.second);
return {field, true};
}
case TypeIndex::Nullable:
{
auto type_nullable = std::static_pointer_cast<const DataTypeNullable>(type);
return {getIcebergType(type_nullable->getNestedType(), iter).first, false};
return {getIcebergType(type_nullable->getNestedType(), iter, format_version).first, false};
}
case TypeIndex::Variant:
{
Expand Down Expand Up @@ -683,12 +705,16 @@ Poco::Dynamic::Var getAvroType(DataTypePtr type, Int32 field_id)
}
case TypeIndex::DateTime64:
{
if (getDecimalScale(*type) != 6)
/// Iceberg `timestamp` is microseconds; `timestamp_ns` is nanoseconds (v3).
/// contrib Avro has no `timestamp-nanos` enum, so it stores the logical type
/// as `NONE` and the value as `long`; other engines still see the annotation.
const auto scale = getDecimalScale(*type);
if (scale != 6 && scale != 9)
throw Exception(ErrorCodes::BAD_ARGUMENTS, "Unsupported type for iceberg {}", type->getName());

Poco::JSON::Object::Ptr timestamp_type = new Poco::JSON::Object;
timestamp_type->set("type", "long");
timestamp_type->set("logicalType", "timestamp-micros");
timestamp_type->set("logicalType", scale == 9 ? "timestamp-nanos" : "timestamp-micros");
timestamp_type->set("adjust-to-utc", assert_cast<const DataTypeDateTime64 &>(*type).hasExplicitTimeZone());
return timestamp_type;
}
Expand Down Expand Up @@ -1086,7 +1112,7 @@ std::pair<Poco::JSON::Object::Ptr, String> createEmptyMetadataFile(
Poco::JSON::Object::Ptr field = new Poco::JSON::Object;
field->set(Iceberg::f_id, ++iter_for_initial_columns);
field->set(Iceberg::f_name, column.name);
auto type = getIcebergType(column.type, iter);
auto type = getIcebergType(column.type, iter, format_version);
field->set(Iceberg::f_required, type.second);
field->set(Iceberg::f_type, type.first);
column_name_to_source_id[column.name] = iter_for_initial_columns;
Expand Down
5 changes: 4 additions & 1 deletion src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.h
Original file line number Diff line number Diff line change
Expand Up @@ -92,7 +92,10 @@ Poco::JSON::Object::Ptr getMetadataJSONObject(
const std::optional<String> & table_uuid);


std::pair<Poco::Dynamic::Var, bool> getIcebergType(DataTypePtr type, Int32 & iter);
/// Maps a ClickHouse type to an Iceberg JSON type. `timestamp_ns` / `timestamptz_ns`
/// require `format_version` >= 3; lower versions throw instead of writing v3 types
/// into v1/v2 metadata.
std::pair<Poco::Dynamic::Var, bool> getIcebergType(DataTypePtr type, Int32 & iter, UInt64 format_version);
Poco::Dynamic::Var getAvroType(DataTypePtr type, Int32 field_id);
Poco::Dynamic::Var getAvroLogicalType(DataTypePtr type);

Expand Down
Loading
Loading