Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions src/Storages/ObjectStorage/DataLakes/Iceberg/Constant.h
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@ DEFINE_ICEBERG_FIELD(logicalType); /// this field has a camelCase name
DEFINE_ICEBERG_FIELD(transform);
DEFINE_ICEBERG_FIELD(direction);

DEFINE_ICEBERG_FIELD(unknown);
DEFINE_ICEBERG_FIELD(uuid);
DEFINE_ICEBERG_FIELD(value);
DEFINE_ICEBERG_FIELD(manifest_length);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -58,8 +58,19 @@ void DataFileStatistics::update(const Chunk & chunk)
}
}

void DataFileStatistics::excludeColumns(std::vector<bool> excluded_)
{
chassert(excluded_.size() == field_ids.size());
excluded = std::move(excluded_);
}

void DataFileStatistics::merge(const DataFileStatistics & other)
{
if (excluded.size() < other.excluded.size())
excluded.resize(other.excluded.size(), false);
for (size_t i = 0; i < other.excluded.size(); ++i)
excluded[i] = excluded[i] || other.excluded[i];

if (other.column_sizes.empty())
return;

Expand Down Expand Up @@ -94,7 +105,8 @@ std::vector<std::pair<size_t, size_t>> DataFileStatistics::getColumnSizes() cons
std::vector<std::pair<size_t, size_t>> result;
for (size_t i = 0; i < column_sizes.size(); ++i)
{
result.push_back({field_ids[i], column_sizes[i]});
if (!isExcluded(i))
result.push_back({field_ids[i], column_sizes[i]});
}
return result;
}
Expand All @@ -104,7 +116,8 @@ std::vector<std::pair<size_t, size_t>> DataFileStatistics::getNullCounts() const
std::vector<std::pair<size_t, size_t>> result;
for (size_t i = 0; i < null_counts.size(); ++i)
{
result.push_back({field_ids[i], null_counts[i]});
if (!isExcluded(i))
result.push_back({field_ids[i], null_counts[i]});
}
return result;
}
Expand All @@ -115,7 +128,8 @@ std::vector<std::pair<size_t, Field>> DataFileStatistics::getLowerBounds() const
std::vector<std::pair<size_t, Field>> result;
for (size_t i = 0; i < ranges.size(); ++i)
{
result.push_back({field_ids[i], ranges[i].left});
if (!isExcluded(i))
result.push_back({field_ids[i], ranges[i].left});
}
return result;
}
Expand All @@ -125,7 +139,8 @@ std::vector<std::pair<size_t, Field>> DataFileStatistics::getUpperBounds() const
std::vector<std::pair<size_t, Field>> result;
for (size_t i = 0; i < ranges.size(); ++i)
{
result.push_back({field_ids[i], ranges[i].right});
if (!isExcluded(i))
result.push_back({field_ids[i], ranges[i].right});
}
return result;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,10 @@ class DataFileStatistics
void update(const Chunk & chunk);
void merge(const DataFileStatistics & other);

/// Marks schema positions whose column is not stored in the data file (e.g. a column of only the
/// Iceberg `unknown` type): the getters omit them, so no statistics describe a column that was not written.
void excludeColumns(std::vector<bool> excluded_);

std::vector<std::pair<size_t, size_t>> getColumnSizes() const;
std::vector<std::pair<size_t, size_t>> getNullCounts() const;
std::vector<std::pair<size_t, Field>> getLowerBounds() const;
Expand All @@ -35,8 +39,10 @@ class DataFileStatistics
const std::vector<Int64> & getFieldIds() const { return field_ids; }
private:
static Range uniteRanges(const Range & left, const Range & right);
bool isExcluded(size_t i) const { return i < excluded.size() && excluded[i]; }

std::vector<Int64> field_ids;
std::vector<bool> excluded;
std::vector<Int64> column_sizes;
std::vector<Int64> null_counts;
std::vector<Range> ranges;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -726,6 +726,7 @@ void MetadataGenerator::generateAddColumnMetadata(const String & column_name, Da
{
if (!type->isNullable())
throw Exception(ErrorCodes::BAD_ARGUMENTS, "Iceberg spec doesn't allow to add non-nullable columns");
Iceberg::checkUnknownTypeAllowed(column_name, type, metadata_object->getValue<Int64>(Iceberg::f_format_version));
const auto next_schema_id = getNextSchemaId(metadata_object);

auto current_schema = deepCopy(getCurrentSchema());
Expand Down Expand Up @@ -756,6 +757,7 @@ void MetadataGenerator::generateAddColumnMetadata(const String & column_name, Da

bool MetadataGenerator::generateModifyColumnMetadata(const String & column_name, DataTypePtr type, ContextPtr context)
{
Iceberg::checkUnknownTypeAllowed(column_name, type, metadata_object->getValue<Int64>(Iceberg::f_format_version));
auto current_schema = getCurrentSchema();

auto last_column_id = metadata_object->getValue<Int32>(Iceberg::f_last_column_id);
Expand Down
241 changes: 239 additions & 2 deletions src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.cpp
Original file line number Diff line number Diff line change
@@ -1,5 +1,15 @@
#include <Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.h>

#include <Columns/ColumnArray.h>
#include <Columns/ColumnMap.h>
#include <Columns/ColumnNullable.h>
#include <Columns/ColumnSparse.h>
#include <Columns/ColumnTuple.h>
#include <DataTypes/DataTypeArray.h>
#include <DataTypes/DataTypeMap.h>
#include <DataTypes/DataTypeNothing.h>
#include <DataTypes/DataTypeNullable.h>
#include <DataTypes/DataTypeTuple.h>
#include <Formats/FormatFactory.h>
#include <Formats/FormatFilterInfo.h>
#include <Processors/Formats/IOutputFormat.h>
Expand All @@ -10,6 +20,147 @@
namespace DB
{

namespace ErrorCodes
{
extern const int LOGICAL_ERROR;
extern const int NOT_IMPLEMENTED;
}

namespace Iceberg
{

DataTypePtr stripNothing(const DataTypePtr & type)
{
if (isNothing(type))
return nullptr;

if (const auto * nullable_type = typeid_cast<const DataTypeNullable *>(type.get()))
{
const auto & nested = nullable_type->getNestedType();
auto stripped_nested = stripNothing(nested);
if (!stripped_nested)
return nullptr;
return stripped_nested == nested ? type : makeNullable(stripped_nested);
}

if (const auto * tuple_type = typeid_cast<const DataTypeTuple *>(type.get()))
{
const auto & elements = tuple_type->getElements();
DataTypes kept_elements;
Strings kept_names;
bool changed = false;
for (size_t i = 0; i < elements.size(); ++i)
{
auto stripped_element = stripNothing(elements[i]);
if (!stripped_element)
{
changed = true;
continue;
}
changed |= stripped_element != elements[i];
kept_elements.push_back(stripped_element);
kept_names.push_back(tuple_type->getNameByPosition(i + 1));
}
if (!changed)
return type;
if (kept_elements.empty())
return nullptr;
if (tuple_type->hasExplicitNames())
return std::make_shared<DataTypeTuple>(kept_elements, kept_names);
return std::make_shared<DataTypeTuple>(kept_elements);
}

if (const auto * array_type = typeid_cast<const DataTypeArray *>(type.get()))
{
const auto & nested = array_type->getNestedType();
auto stripped_nested = stripNothing(nested);
if (!stripped_nested)
return nullptr;
return stripped_nested == nested ? type : std::make_shared<DataTypeArray>(stripped_nested);
}

if (const auto * map_type = typeid_cast<const DataTypeMap *>(type.get()))
{
const auto & key = map_type->getKeyType();
const auto & value = map_type->getValueType();
auto stripped_key = stripNothing(key);
auto stripped_value = stripNothing(value);
if (!stripped_key || !stripped_value)
return nullptr;
if (stripped_key == key && stripped_value == value)
return type;
return std::make_shared<DataTypeMap>(stripped_key, stripped_value);
}

return type;
}

namespace
{

ColumnPtr stripNothingColumnImpl(const ColumnPtr & column, const DataTypePtr & type)
{
auto stripped_type = stripNothing(type);
if (!stripped_type)
throw Exception(ErrorCodes::LOGICAL_ERROR, "Column of type {} has nothing left after stripping Nothing", type->getName());
if (stripped_type == type)
return column;

auto full_column = removeSpecialRepresentations(column->convertToFullColumnIfConst());

if (const auto * nullable_type = typeid_cast<const DataTypeNullable *>(type.get()))
{
const auto & nullable_column = assert_cast<const ColumnNullable &>(*full_column);
return ColumnNullable::create(
stripNothingColumnImpl(nullable_column.getNestedColumnPtr(), nullable_type->getNestedType()),
nullable_column.getNullMapColumnPtr());
}

if (const auto * tuple_type = typeid_cast<const DataTypeTuple *>(type.get()))
{
const auto & tuple_column = assert_cast<const ColumnTuple &>(*full_column);
const auto & elements = tuple_type->getElements();
Columns kept_columns;
for (size_t i = 0; i < elements.size(); ++i)
{
if (stripNothing(elements[i]))
kept_columns.push_back(stripNothingColumnImpl(tuple_column.getColumnPtr(i), elements[i]));
}
return ColumnTuple::create(kept_columns);
}

if (const auto * array_type = typeid_cast<const DataTypeArray *>(type.get()))
{
const auto & array_column = assert_cast<const ColumnArray &>(*full_column);
return ColumnArray::create(
stripNothingColumnImpl(array_column.getDataPtr(), array_type->getNestedType()),
array_column.getOffsetsPtr());
}

if (const auto * map_type = typeid_cast<const DataTypeMap *>(type.get()))
{
const auto & map_column = assert_cast<const ColumnMap &>(*full_column);
const auto & key_value = map_column.getNestedData();
return ColumnMap::create(
stripNothingColumnImpl(key_value.getColumnPtr(0), map_type->getKeyType()),
stripNothingColumnImpl(key_value.getColumnPtr(1), map_type->getValueType()),
map_column.getNestedColumn().getOffsetsPtr());
}

throw Exception(ErrorCodes::LOGICAL_ERROR, "Unexpected type {} while stripping Nothing from a column", type->getName());
}

}

ColumnPtr stripNothingColumn(const ColumnPtr & column, const DataTypePtr & original_type, const DataTypePtr & stripped_type)
{
auto result = stripNothingColumnImpl(column, original_type);
chassert(stripped_type && stripped_type->equals(*stripNothing(original_type)));
return result;
}

}

#if USE_AVRO

MultipleFileWriter::MultipleFileWriter(
Expand Down Expand Up @@ -39,6 +190,32 @@ MultipleFileWriter::MultipleFileWriter(
, new_file_path_callback(std::move(new_file_path_callback_))
{
column_mapper->setStorageColumnEncoding(Iceberg::IcebergSchemaProcessor::traverseSchema(schema_));

written_column_types.reserve(sample_block->columns());
for (const auto & column : *sample_block)
{
written_column_types.push_back(Iceberg::stripNothing(column.type));
has_nothing_leaves |= written_column_types.back() != column.type;
}

if (!has_nothing_leaves)
{
filtered_sample_block = sample_block;
}
else
{
aggregate_stats.excludeColumns(getStatisticsExcludedColumns());

Block filtered;
for (size_t i = 0; i < sample_block->columns(); ++i)
{
if (!written_column_types[i])
continue;
const auto & column = sample_block->getByPosition(i);
filtered.insert({written_column_types[i]->createColumn(), written_column_types[i], column.name});
}
filtered_sample_block = std::make_shared<const Block>(std::move(filtered));
}
}

void MultipleFileWriter::startNewFile()
Expand All @@ -49,6 +226,8 @@ void MultipleFileWriter::startNewFile()
}

current_file_stats = std::make_shared<DataFileStatistics>(schema);
if (has_nothing_leaves)
current_file_stats->excludeColumns(getStatisticsExcludedColumns());
current_file_num_rows = 0;
current_file_num_bytes = 0;
auto metadata_path = filename_generator.generateDataFileName();
Expand All @@ -69,16 +248,74 @@ void MultipleFileWriter::startNewFile()
}
FormatFilterInfoPtr format_filter_info = std::make_shared<FormatFilterInfo>(nullptr, context, column_mapper, nullptr, nullptr);
output_format = FormatFactory::instance().getOutputFormatParallelIfPossible(
write_format, *buffer, *sample_block, context, format_settings, format_filter_info);
write_format, *buffer, *filtered_sample_block, context, format_settings, format_filter_info);
}

std::vector<bool> MultipleFileWriter::getStatisticsExcludedColumns() const
{
std::vector<bool> excluded(written_column_types.size());
for (size_t i = 0; i < written_column_types.size(); ++i)
excluded[i] = !written_column_types[i];
return excluded;
}

Columns MultipleFileWriter::filterColumns(const Columns & columns) const
{
Columns filtered_columns;
filtered_columns.reserve(filtered_sample_block->columns());
for (size_t i = 0; i < columns.size(); ++i)
{
const auto & original_type = sample_block->getByPosition(i).type;
const auto & written_type = written_column_types[i];
if (written_type == original_type)
{
filtered_columns.push_back(columns[i]);
}
else if (written_type)
{
filtered_columns.push_back(Iceberg::stripNothingColumn(columns[i], original_type, written_type));
}
else
{
for (size_t row = 0; row < columns[i]->size(); ++row)
{
if (!columns[i]->isDefaultAt(row))
throw Exception(
ErrorCodes::NOT_IMPLEMENTED,
"Cannot write column '{}' of type {} into an Iceberg data file: row {} holds a non-default value, "
"but the column contains an Iceberg `unknown` element that no data file format can store, so the "
"value (for example, list elements or map keys) would be lost. Only NULL, empty lists and empty "
"maps can be inserted into such a column",
sample_block->getByPosition(i).name, original_type->getName(), row);
}
}
}
return filtered_columns;
}

void MultipleFileWriter::consume(const Chunk & chunk)
{
/// Validate before starting a file, so a rejected chunk leaves no empty data file behind.
if (filtered_sample_block->columns() == 0 && sample_block->columns() > 0)
throw Exception(
ErrorCodes::NOT_IMPLEMENTED,
"Cannot write an Iceberg data file: every column contains only the Iceberg `unknown` type, "
"which no data file format can store");

std::optional<Columns> filtered_columns;
if (has_nothing_leaves)
filtered_columns = filterColumns(chunk.getColumns());

if (!current_file_num_rows || *current_file_num_rows >= max_data_file_num_rows || *current_file_num_bytes >= max_data_file_num_bytes)
{
startNewFile();
}
output_format->write(sample_block->cloneWithColumns(chunk.getColumns()));

if (filtered_columns)
output_format->write(filtered_sample_block->cloneWithColumns(std::move(*filtered_columns)));
else
output_format->write(sample_block->cloneWithColumns(chunk.getColumns()));

output_format->flush();
*current_file_num_rows += chunk.getNumRows();
*current_file_num_bytes += chunk.bytes();
Expand Down
Loading
Loading