Skip to content

Commit 911f0dc

Browse files
authored
Merge pull request #2125 from Altinity/feature/antalya-26.6/pr-1655
Antalya 26.6: Implement TRUNCATE TABLE for Iceberg Engine (REST …
2 parents be5f3f6 + dfeeb88 commit 911f0dc

8 files changed

Lines changed: 257 additions & 10 deletions

File tree

‎src/Storages/ObjectStorage/DataLakes/IDataLakeMetadata.h‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -244,6 +244,12 @@ class IDataLakeMetadata : boost::noncopyable
244244
throwNotImplemented(fmt::format("EXECUTE {}", command_name));
245245
}
246246

247+
virtual bool supportsTruncate() const { return false; }
248+
virtual void truncate(ContextPtr /*context*/, std::shared_ptr<DataLake::ICatalog> /*catalog*/, const StorageID & /*storage_id*/)
249+
{
250+
throwNotImplemented("truncate");
251+
}
252+
247253
virtual void drop(ContextPtr) { }
248254

249255
virtual ObjectStorageType getObjectStorageType() const { return ObjectStorageType::None; }

‎src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.cpp‎

Lines changed: 85 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -73,6 +73,7 @@
7373
#include <Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.h>
7474
#include <Storages/ObjectStorage/DataLakes/Iceberg/IcebergTableStateSnapshot.h>
7575
#include <Storages/ObjectStorage/DataLakes/Iceberg/IcebergWrites.h>
76+
#include <Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.h>
7677
#include <Storages/ObjectStorage/DataLakes/Iceberg/ManifestFile.h>
7778
#include <Storages/ObjectStorage/DataLakes/Iceberg/ManifestFilesPruning.h>
7879
#include <Storages/ObjectStorage/DataLakes/Iceberg/Mutations.h>
@@ -115,6 +116,7 @@ extern const int ICEBERG_SPECIFICATION_VIOLATION;
115116
extern const int S3_ERROR;
116117
extern const int TABLE_ALREADY_EXISTS;
117118
extern const int SUPPORT_IS_DISABLED;
119+
extern const int INCORRECT_DATA;
118120
}
119121

120122
namespace Setting
@@ -277,15 +279,17 @@ void IcebergMetadata::backgroundMetadataPrefetcherThread()
277279
/// first, we fetch the latest metadata version and cache it;
278280
/// as a part of the same method, we download metadata.json of the latest metadata version
279281
/// and after parsing it, we fetch manifest lists, parse and cache them
280-
auto ctx = Context::getGlobalContextInstance()->getBackgroundContext();
282+
auto ctx = Context::createCopy(Context::getGlobalContextInstance());
281283
auto [actual_data_snapshot, actual_table_state_snapshot] = getRelevantState(ctx, true);
282284
if (actual_data_snapshot)
283285
{
284286
for (const auto & entry : actual_data_snapshot->manifest_list_entries)
285287
{
286288
/// second, we fetch, parse and cache each manifest file
287-
auto manifest_file_ptr = getManifestFileEntriesHandle(
288-
object_storage, persistent_components, ctx, log, entry, actual_table_state_snapshot.schema_id);
289+
auto manifest_file_ptr = Iceberg::getManifestFile(
290+
object_storage, persistent_components, ctx, log,
291+
entry.manifest_file_path,
292+
entry.manifest_file_byte_size);
289293
}
290294
}
291295

@@ -611,6 +615,84 @@ void IcebergMetadata::mutate(
611615
catalog);
612616
}
613617

618+
void IcebergMetadata::truncate(ContextPtr context, std::shared_ptr<DataLake::ICatalog> catalog, const StorageID & storage_id)
619+
{
620+
if (!context->getSettingsRef()[Setting::allow_insert_into_iceberg].value)
621+
throw Exception(
622+
ErrorCodes::SUPPORT_IS_DISABLED,
623+
"Iceberg truncate requires the allow_insert_into_iceberg setting to be enabled.");
624+
625+
auto [actual_data_snapshot, actual_table_state_snapshot] = getRelevantState(context);
626+
auto metadata_object = getMetadataJSONObject(
627+
actual_table_state_snapshot.metadata_file_path,
628+
object_storage,
629+
persistent_components.metadata_cache,
630+
context,
631+
log,
632+
persistent_components.metadata_compression_method,
633+
persistent_components.table_uuid);
634+
635+
// Use -1 as the Iceberg spec sentinel for "no parent snapshot"
636+
// (distinct from snapshot ID 0 which is a valid snapshot).
637+
Int64 parent_snapshot_id = actual_table_state_snapshot.snapshot_id.value_or(-1);
638+
639+
// On antalya-26.6 all metadata paths flow through the table's IcebergPathResolver:
640+
// FileNamesGenerator produces IcebergPathFromMetadata values (relative to the table
641+
// location), and resolver.resolve() / resolveForCatalog() turn those into the storage
642+
// path for I/O and the fully-qualified path the catalog expects. This mirrors the
643+
// write path in IcebergStorageSink (see IcebergWrites.cpp) so transactional (REST) and
644+
// non-transactional catalogs are handled uniformly.
645+
const auto & resolver = persistent_components.path_resolver;
646+
647+
bool is_transactional = (catalog != nullptr && catalog->isTransactional());
648+
649+
FileNamesGenerator filename_generator(
650+
resolver.getTableLocation(),
651+
is_transactional,
652+
persistent_components.metadata_compression_method,
653+
write_format);
654+
655+
Int32 new_metadata_version = actual_table_state_snapshot.metadata_version + 1;
656+
filename_generator.setVersion(new_metadata_version);
657+
658+
auto metadata_info = filename_generator.generateMetadataPathWithInfo();
659+
660+
auto [new_snapshot, manifest_list_path] = MetadataGenerator(metadata_object).generateNextMetadata(
661+
filename_generator, metadata_info.path, parent_snapshot_id,
662+
/* added_files */ 0, /* added_records */ 0, /* added_files_size */ 0,
663+
/* num_partitions */ 0, /* added_delete_files */ 0, /* num_deleted_rows */ 0,
664+
std::nullopt, std::nullopt, /*is_truncate=*/true);
665+
666+
auto storage_manifest_list_name = resolver.resolve(manifest_list_path);
667+
668+
auto write_settings = context->getWriteSettings();
669+
auto buf = object_storage->writeObject(
670+
StoredObject(storage_manifest_list_name),
671+
WriteMode::Rewrite, std::nullopt,
672+
DBMS_DEFAULT_BUFFER_SIZE, write_settings);
673+
674+
// Truncate writes a metadata-only overwrite snapshot: an empty manifest list
675+
// (no manifest entries, no sizes) that supersedes all previous snapshots.
676+
generateManifestList(resolver, metadata_object, object_storage,
677+
context, {}, new_snapshot, {}, *buf, Iceberg::FileContentType::DATA, /*use_previous_snapshots=*/false);
678+
buf->finalize();
679+
680+
String metadata_content = dumpMetadataObjectToString(metadata_object);
681+
writeMessageToFile(metadata_content, resolver.resolve(metadata_info.path), object_storage,
682+
context, "*", "", persistent_components.metadata_compression_method);
683+
684+
if (catalog)
685+
{
686+
String catalog_filename = resolver.resolveForCatalog(metadata_info.path);
687+
688+
const auto & [namespace_name, table_name] = DataLake::parseTableName(storage_id.getTableName());
689+
if (!catalog->updateMetadata(namespace_name, table_name, catalog_filename, new_snapshot))
690+
throw Exception(ErrorCodes::INCORRECT_DATA,
691+
"Failed to commit Iceberg truncate update to catalog.");
692+
}
693+
}
694+
695+
614696
void IcebergMetadata::checkMutationIsPossible(const MutationCommands & commands)
615697
{
616698
for (const auto & command : commands)

‎src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergMetadata.h‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -104,6 +104,8 @@ class IcebergMetadata : public IDataLakeMetadata
104104
bool supportsUpdate() const override { return true; }
105105
bool supportsWrites() const override { return true; }
106106
bool supportsParallelInsert() const override { return true; }
107+
bool supportsTruncate() const override { return true; }
108+
void truncate(ContextPtr context, std::shared_ptr<DataLake::ICatalog> catalog, const StorageID & storage_id) override;
107109

108110
IcebergHistory getHistory(ContextPtr local_context) const;
109111

‎src/Storages/ObjectStorage/DataLakes/Iceberg/IcebergWrites.cpp‎

Lines changed: 57 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -539,6 +539,31 @@ void generateManifestFile(
539539
writer.close();
540540
}
541541

542+
// Avro uses zigzag encoding for integers to efficiently represent small negative
543+
// numbers. Positive n maps to 2n, negative n maps to 2(-n)-1, keeping small
544+
// magnitudes compact regardless of sign. The value is then serialized as a
545+
// variable-length base-128 integer (little-endian), where the high bit of each
546+
// byte signals whether more bytes follow.
547+
// See: https://avro.apache.org/docs/1.11.1/specification/#binary-encoding
548+
static void writeAvroLong(WriteBuffer & out, int64_t val)
549+
{
550+
uint64_t n = (static_cast<uint64_t>(val) << 1) ^ static_cast<uint64_t>(val >> 63);
551+
while (n & ~0x7fULL)
552+
{
553+
char c = static_cast<char>((n & 0x7f) | 0x80);
554+
out.write(&c, 1);
555+
n >>= 7;
556+
}
557+
char c = static_cast<char>(n);
558+
out.write(&c, 1);
559+
}
560+
561+
static void writeAvroBytes(WriteBuffer & out, const String & s)
562+
{
563+
writeAvroLong(out, static_cast<int64_t>(s.size()));
564+
out.write(s.data(), s.size());
565+
}
566+
542567
void generateManifestList(
543568
const Iceberg::IcebergPathResolver & path_resolver,
544569
Poco::JSON::Object::Ptr metadata,
@@ -558,6 +583,38 @@ void generateManifestList(
558583
else
559584
schema_representation = manifest_list_v2_schema;
560585

586+
// For empty manifest list (e.g. TRUNCATE), write a valid Avro container
587+
// file manually so we can embed the full schema JSON with field-ids intact,
588+
// without triggering the DataFileWriter constructor's eager writeHeader()
589+
// which commits encoder state before we can override avro.schema.
590+
if (manifest_entry_names.empty() && !use_previous_snapshots)
591+
{
592+
// For an empty manifest list (e.g. after TRUNCATE), we write a minimal valid
593+
// Avro Object Container File manually rather than using avro::DataFileWriter.
594+
// The reason: DataFileWriter calls writeHeader() eagerly in its constructor,
595+
// committing the binary encoder state. Post-construction setMetadata() calls
596+
// corrupt StreamWriter::next_ causing a NULL dereference on close(). Writing
597+
// the OCF header directly ensures the full schema JSON (with Iceberg field-ids)
598+
// is embedded intact — the Avro C++ library strips unknown field properties
599+
// like field-id during schema node serialization.
600+
// Avro OCF format: [magic(4)] [metadata_map] [sync_marker(16)] [no data blocks]
601+
buf.write("Obj\x01", 4);
602+
603+
writeAvroLong(buf, 2); // 2 metadata entries
604+
writeAvroBytes(buf, "avro.codec");
605+
writeAvroBytes(buf, "null");
606+
writeAvroBytes(buf, "avro.schema");
607+
writeAvroBytes(buf, schema_representation); // full JSON with field-ids intact
608+
609+
writeAvroLong(buf, 0); // end of metadata map
610+
611+
static const char sync_marker[16] = {};
612+
buf.write(sync_marker, 16);
613+
614+
buf.finalize();
615+
return;
616+
}
617+
561618
auto schema = avro::compileJsonSchemaFromString(schema_representation); // NOLINT
562619

563620
auto adapter = std::make_unique<OutputStreamWriteBufferAdapter>(buf);

‎src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.cpp‎

Lines changed: 18 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -121,7 +121,8 @@ MetadataGenerator::NextMetadataResult MetadataGenerator::generateNextMetadata(
121121
Int64 added_delete_files,
122122
Int64 num_deleted_rows,
123123
std::optional<Int64> user_defined_snapshot_id,
124-
std::optional<Int64> user_defined_timestamp)
124+
std::optional<Int64> user_defined_timestamp,
125+
bool is_truncate)
125126
{
126127
int format_version = metadata_object->getValue<Int32>(Iceberg::f_format_version);
127128
Poco::JSON::Object::Ptr new_snapshot = new Poco::JSON::Object;
@@ -145,7 +146,16 @@ MetadataGenerator::NextMetadataResult MetadataGenerator::generateNextMetadata(
145146

146147
auto parent_snapshot = getParentSnapshot(parent_snapshot_id);
147148
Poco::JSON::Object::Ptr summary = new Poco::JSON::Object;
148-
if (num_deleted_rows == 0)
149+
if (is_truncate)
150+
{
151+
summary->set(Iceberg::f_operation, Iceberg::f_overwrite);
152+
Int32 prev_total_records = parent_snapshot && parent_snapshot->has(Iceberg::f_summary) && parent_snapshot->getObject(Iceberg::f_summary)->has(Iceberg::f_total_records) ? std::stoi(parent_snapshot->getObject(Iceberg::f_summary)->getValue<String>(Iceberg::f_total_records)) : 0;
153+
Int32 prev_total_data_files = parent_snapshot && parent_snapshot->has(Iceberg::f_summary) && parent_snapshot->getObject(Iceberg::f_summary)->has(Iceberg::f_total_data_files) ? std::stoi(parent_snapshot->getObject(Iceberg::f_summary)->getValue<String>(Iceberg::f_total_data_files)) : 0;
154+
155+
summary->set(Iceberg::f_deleted_records, std::to_string(prev_total_records));
156+
summary->set(Iceberg::f_deleted_data_files, std::to_string(prev_total_data_files));
157+
}
158+
else if (num_deleted_rows == 0)
149159
{
150160
summary->set(Iceberg::f_operation, Iceberg::f_append);
151161
summary->set(Iceberg::f_added_data_files, std::to_string(added_files));
@@ -165,7 +175,12 @@ MetadataGenerator::NextMetadataResult MetadataGenerator::generateNextMetadata(
165175

166176
auto sum_with_parent_snapshot = [&](const char * field_name, Int64 snapshot_value)
167177
{
168-
Int64 prev_value = parent_snapshot ? parse<Int64>(parent_snapshot->getObject(Iceberg::f_summary)->getValue<String>(field_name)) : 0;
178+
if (is_truncate)
179+
{
180+
summary->set(field_name, std::to_string(0));
181+
return;
182+
}
183+
Int64 prev_value = parent_snapshot && parent_snapshot->has(Iceberg::f_summary) && parent_snapshot->getObject(Iceberg::f_summary)->has(field_name) ? parse<Int64>(parent_snapshot->getObject(Iceberg::f_summary)->getValue<String>(field_name)) : 0;
169184
summary->set(field_name, std::to_string(prev_value + snapshot_value));
170185
};
171186

‎src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.h‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -39,7 +39,8 @@ class MetadataGenerator
3939
Int64 added_delete_files,
4040
Int64 num_deleted_rows,
4141
std::optional<Int64> user_defined_snapshot_id = std::nullopt,
42-
std::optional<Int64> user_defined_timestamp = std::nullopt);
42+
std::optional<Int64> user_defined_timestamp = std::nullopt,
43+
bool is_truncate = false);
4344

4445
void generateAddColumnMetadata(const String & column_name, DataTypePtr type);
4546
void generateDropColumnMetadata(const String & column_name);

‎src/Storages/ObjectStorage/StorageObjectStorage.cpp‎

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -651,7 +651,7 @@ bool StorageObjectStorage::optimize(
651651
void StorageObjectStorage::truncate(
652652
const ASTPtr & /* query */,
653653
const StorageMetadataPtr & /* metadata_snapshot */,
654-
ContextPtr /* context */,
654+
ContextPtr local_context,
655655
TableExclusiveLockHolder & /* table_holder */)
656656
{
657657
const auto path = configuration->getRawPath();
@@ -665,8 +665,12 @@ void StorageObjectStorage::truncate(
665665

666666
if (configuration->isDataLakeConfiguration())
667667
{
668-
throw Exception(ErrorCodes::NOT_IMPLEMENTED,
669-
"Truncate is not supported for data lake engine");
668+
auto * data_lake_metadata = getExternalMetadata(local_context);
669+
if (!data_lake_metadata || !data_lake_metadata->supportsTruncate())
670+
throw Exception(ErrorCodes::NOT_IMPLEMENTED, "Truncate is not supported for this data lake engine");
671+
672+
data_lake_metadata->truncate(local_context, catalog, getStorageID());
673+
return;
670674
}
671675

672676
if (path.hasGlobsIgnorePlaceholders())
Lines changed: 80 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,80 @@
1+
#!/usr/bin/env python3
2+
3+
from pyiceberg.catalog import load_catalog
4+
from helpers.config_cluster import minio_secret_key, minio_access_key
5+
import uuid
6+
import pyarrow as pa
7+
from pyiceberg.schema import Schema, NestedField
8+
from pyiceberg.types import LongType, StringType
9+
from pyiceberg.partitioning import PartitionSpec
10+
11+
CATALOG_NAME = "demo"
12+
13+
def load_catalog_impl(started_cluster):
14+
return load_catalog(
15+
CATALOG_NAME,
16+
**{
17+
"uri": f"http://localhost:{started_cluster.iceberg_rest_catalog_port}",
18+
"type": "rest",
19+
"s3.endpoint": f"http://{started_cluster.minio_ip}:{started_cluster.minio_port}",
20+
"s3.access-key-id": minio_access_key,
21+
"s3.secret-access-key": minio_secret_key,
22+
},
23+
)
24+
25+
26+
def test_iceberg_truncate_restart(started_cluster_iceberg_no_spark):
27+
instance = started_cluster_iceberg_no_spark.instances["node1"]
28+
catalog = load_catalog_impl(started_cluster_iceberg_no_spark)
29+
30+
namespace = f"clickhouse_truncate_restart_{uuid.uuid4().hex}"
31+
catalog.create_namespace(namespace)
32+
33+
schema = Schema(
34+
NestedField(field_id=1, name="id", field_type=LongType(), required=False),
35+
NestedField(field_id=2, name="val", field_type=StringType(), required=False),
36+
)
37+
table_name = "test_truncate_restart"
38+
catalog.create_table(
39+
identifier=f"{namespace}.{table_name}",
40+
schema=schema,
41+
location=f"s3://warehouse-rest/{namespace}.{table_name}",
42+
partition_spec=PartitionSpec(),
43+
)
44+
45+
ch_table_identifier = f"`{namespace}.{table_name}`"
46+
47+
instance.query(f"DROP DATABASE IF EXISTS {namespace}")
48+
instance.query(
49+
f"""
50+
CREATE DATABASE {namespace} ENGINE = DataLakeCatalog('http://rest:8181/v1', 'minio', '{minio_secret_key}')
51+
SETTINGS
52+
catalog_type='rest',
53+
warehouse='demo',
54+
storage_endpoint='http://minio1:9001/warehouse-rest';
55+
""",
56+
settings={"allow_database_iceberg": 1}
57+
)
58+
59+
# 1. Insert initial data and truncate
60+
df = pa.Table.from_pylist([{"id": 1, "val": "A"}, {"id": 2, "val": "B"}])
61+
catalog.load_table(f"{namespace}.{table_name}").append(df)
62+
63+
assert int(instance.query(f"SELECT count() FROM {namespace}.{ch_table_identifier}").strip()) == 2
64+
65+
instance.query(
66+
f"TRUNCATE TABLE {namespace}.{ch_table_identifier}",
67+
settings={"allow_experimental_insert_into_iceberg": 1}
68+
)
69+
assert int(instance.query(f"SELECT count() FROM {namespace}.{ch_table_identifier}").strip()) == 0
70+
71+
# 2. Restart ClickHouse and verify table is still readable (count = 0)
72+
instance.restart_clickhouse()
73+
assert int(instance.query(f"SELECT count() FROM {namespace}.{ch_table_identifier}").strip()) == 0
74+
75+
# 3. Insert new data after restart and verify it's readable
76+
new_df = pa.Table.from_pylist([{"id": 3, "val": "C"}])
77+
catalog.load_table(f"{namespace}.{table_name}").append(new_df)
78+
assert int(instance.query(f"SELECT count() FROM {namespace}.{ch_table_identifier}").strip()) == 1
79+
80+
instance.query(f"DROP DATABASE {namespace}")

0 commit comments

Comments
 (0)