From 48d85558557c65f7e0da100e069dc05cc1725092 Mon Sep 17 00:00:00 2001 From: Aykut Bozkurt Date: Fri, 9 Jan 2026 19:09:46 +0300 Subject: [PATCH 1/3] Refactor IcebergPartitionTransform IcebergPartitionTransform can embed a IcebergPartitionSpecField by removing some individual fields, which makes it less verbose. Signed-off-by: Aykut Bozkurt --- .../src/data_file/data_file_stats.c | 2 +- .../pg_lake/iceberg/api/partitioning.h | 10 ++------- .../src/iceberg/partitioning/partition.c | 2 +- .../iceberg/partitioning/spec_generation.c | 22 +------------------ .../partitioning/partition_spec_catalog.h | 1 - pg_lake_table/src/fdw/partition_transform.c | 14 +++++------- .../fdw/partitioning/partition_by_parser.c | 18 +++++++++++---- 7 files changed, 25 insertions(+), 44 deletions(-) diff --git a/pg_lake_engine/src/data_file/data_file_stats.c b/pg_lake_engine/src/data_file/data_file_stats.c index 1d78c82a7..5e5bbb43a 100644 --- a/pg_lake_engine/src/data_file/data_file_stats.c +++ b/pg_lake_engine/src/data_file/data_file_stats.c @@ -256,7 +256,7 @@ ExtractMinMaxForColumn(Datum map, char *colName, List **names, List **mins, List if (minText != NULL && maxText != NULL) { - *names = lappend(*names, colName); + *names = lappend(*names, pstrdup(colName)); *mins = lappend(*mins, minText); *maxs = lappend(*maxs, maxText); } diff --git a/pg_lake_iceberg/include/pg_lake/iceberg/api/partitioning.h b/pg_lake_iceberg/include/pg_lake/iceberg/api/partitioning.h index aa2cdfc15..535d436c4 100644 --- a/pg_lake_iceberg/include/pg_lake/iceberg/api/partitioning.h +++ b/pg_lake_iceberg/include/pg_lake/iceberg/api/partitioning.h @@ -19,6 +19,7 @@ #include "postgres.h" +#include "pg_lake/iceberg/metadata_spec.h" #include "pg_lake/parquet/field.h" #include "pg_lake/pgduck/type.h" #include "access/attnum.h" @@ -51,14 +52,7 @@ typedef struct IcebergPartitionTransform size_t truncateLen; }; - /* partition field id */ - int32_t partitionFieldId; - - /* _, e.g. a_bucket */ - const char *partitionFieldName; - - /* transform name, e.g. bucket[3] */ - const char *transformName; + IcebergPartitionSpecField *specField; /* source field of the column to which transform applies */ DataFileSchemaField *sourceField; diff --git a/pg_lake_iceberg/src/iceberg/partitioning/partition.c b/pg_lake_iceberg/src/iceberg/partitioning/partition.c index 141f7c974..1798e067b 100644 --- a/pg_lake_iceberg/src/iceberg/partitioning/partition.c +++ b/pg_lake_iceberg/src/iceberg/partitioning/partition.c @@ -225,7 +225,7 @@ FindPartitionTransformById(List *transforms, int32_t partitionFieldId, bool erro { IcebergPartitionTransform *transform = (IcebergPartitionTransform *) lfirst(cell); - if (transform->partitionFieldId == partitionFieldId) + if (transform->specField->field_id == partitionFieldId) return transform; } diff --git a/pg_lake_iceberg/src/iceberg/partitioning/spec_generation.c b/pg_lake_iceberg/src/iceberg/partitioning/spec_generation.c index 83f2f0d11..ccd5762cb 100644 --- a/pg_lake_iceberg/src/iceberg/partitioning/spec_generation.c +++ b/pg_lake_iceberg/src/iceberg/partitioning/spec_generation.c @@ -60,27 +60,7 @@ BuildPartitionSpecFromPartitionTransforms(Oid relationId, List *partitionTransfo { IcebergPartitionTransform *transform = lfirst(transformCell); - IcebergPartitionSpecField *field = palloc0(sizeof(IcebergPartitionSpecField)); - - field->source_id = transform->sourceField->id; - - /* - * We do not support partition transforms on multi columns (v3 - * feature), and to comply with the iceberg spec/reference - * implementation for v2, we still fill the source_ids array. - */ - field->source_ids_length = 1; - field->source_ids = palloc0(sizeof(int) * field->source_ids_length); - field->source_ids[0] = transform->sourceField->id; - - field->field_id = transform->partitionFieldId; - - field->name = pstrdup(transform->partitionFieldName); - field->name_length = strlen(transform->partitionFieldName); - field->transform = pstrdup(transform->transformName); - field->transform_length = strlen(transform->transformName); - - spec->fields[fieldIndex] = *field; + spec->fields[fieldIndex] = *(transform->specField); fieldIndex++; } diff --git a/pg_lake_table/include/pg_lake/partitioning/partition_spec_catalog.h b/pg_lake_table/include/pg_lake/partitioning/partition_spec_catalog.h index aefc6358d..fe2d10ebc 100644 --- a/pg_lake_table/include/pg_lake/partitioning/partition_spec_catalog.h +++ b/pg_lake_table/include/pg_lake/partitioning/partition_spec_catalog.h @@ -43,7 +43,6 @@ typedef struct IcebergPartitionSpecHashEntry extern void UpdateDefaultPartitionSpecId(Oid relationId, int specId); extern void InsertPartitionSpecAndPartitionFields(Oid relationId, IcebergPartitionSpec * spec); extern int GetLargestSpecId(Oid relationId); -extern List *GetAllIcebergPartitionSpecIds(Oid relationId); extern PGDLLEXPORT int GetCurrentSpecId(Oid relationId); extern int GetLargestPartitionFieldId(Oid relationId); extern IcebergPartitionSpecField * GetIcebergPartitionFieldFromCatalog(Oid relationId, int fieldId); diff --git a/pg_lake_table/src/fdw/partition_transform.c b/pg_lake_table/src/fdw/partition_transform.c index 7d4430e72..4ece3ee4e 100644 --- a/pg_lake_table/src/fdw/partition_transform.c +++ b/pg_lake_table/src/fdw/partition_transform.c @@ -162,7 +162,7 @@ PartitionTransformsEqual(IcebergPartitionSpec * spec, List *partitionTransforms) * Iceberg does here: * https://github.com/apache/iceberg/blob/8b55ac834015ce664f879ecfe1e80a941a994420/api/src/main/java/org/apache/iceberg/PartitionSpec.java#L239-L259 */ - if (strcasecmp(specField->name, transform->partitionFieldName) != 0) + if (strcasecmp(specField->name, transform->specField->name) != 0) { return false; } @@ -251,9 +251,7 @@ GetPartitionTransformFromSpecField(Oid relationId, IcebergPartitionSpecField * s { IcebergPartitionTransform *transform = palloc0(sizeof(IcebergPartitionTransform)); - transform->partitionFieldId = specField->field_id; - transform->partitionFieldName = pstrdup(specField->name); - transform->transformName = pstrdup(specField->transform); + transform->specField = specField; transform->attnum = GetAttributeForFieldId(relationId, specField->source_id); @@ -274,7 +272,7 @@ GetPartitionTransformFromSpecField(Oid relationId, IcebergPartitionSpecField * s } /* parse transform name */ - ParseTransformName(transform->transformName, + ParseTransformName(transform->specField->transform, &transform->type, &transform->bucketCount, &transform->truncateLen); @@ -413,8 +411,8 @@ ApplyPartitionTransformToTuple(IcebergPartitionTransform * transform, TupleTable { PartitionField *field = palloc0(sizeof(PartitionField)); - field->field_name = pstrdup(transform->partitionFieldName); - field->field_id = transform->partitionFieldId; + field->field_name = pstrdup(transform->specField->name); + field->field_id = transform->specField->field_id; bool isNull = false; Datum columnValue = slot_getattr(slot, transform->attnum, &isNull); @@ -453,7 +451,7 @@ ApplyPartitionTransformToTuple(IcebergPartitionTransform * transform, TupleTable ereport(ERROR, (errcode(ERRCODE_FEATURE_NOT_SUPPORTED), errmsg("applying transform %s is not yet support ", - transform->transformName))); + transform->specField->transform))); } field->value_type = GetTransformResultAvroType(transform); diff --git a/pg_lake_table/src/fdw/partitioning/partition_by_parser.c b/pg_lake_table/src/fdw/partitioning/partition_by_parser.c index 25cc16ac2..697c487dc 100644 --- a/pg_lake_table/src/fdw/partitioning/partition_by_parser.c +++ b/pg_lake_table/src/fdw/partitioning/partition_by_parser.c @@ -529,14 +529,24 @@ AnalyzeIcebergTablePartitionBy(Oid relationId, List *transforms) transform->sourceField = sourceField; + transform->specField = palloc0(sizeof(IcebergPartitionSpecField)); + /* set transform name */ - transform->transformName = GenerateTransformName(transform); + transform->specField->transform = GenerateTransformName(transform); + transform->specField->transform_length = strlen(transform->specField->transform); /* set partition field name */ - transform->partitionFieldName = GeneratePartitionFieldName(transform, relationId); + transform->specField->name = GeneratePartitionFieldName(transform, relationId); + transform->specField->name_length = strlen(transform->specField->name); /* set partition field id */ - transform->partitionFieldId = ++largestPartitionFieldId; + transform->specField->field_id = ++largestPartitionFieldId; + + /* set source field id */ + transform->specField->source_id = sourceField->id; + transform->specField->source_ids_length = 1; + transform->specField->source_ids = palloc0(sizeof(int) * transform->specField->source_ids_length); + transform->specField->source_ids[0] = transform->specField->source_id; /* 3) Check column type compatibility. */ EnsureValidTypeForTransform(transform->type, transform->pgType.postgresTypeOid); @@ -822,7 +832,7 @@ EnsureNoDuplicateTransforms(List *transforms) ereport(ERROR, (errcode(ERRCODE_DUPLICATE_OBJECT), errmsg("\"%s\" transform on column \"%s\" appears multiple times in partition spec", - transform->transformName, transform->columnName))); + transform->specField->transform, transform->columnName))); } } } From a90584748923e87a997eecb252888512b75764a4 Mon Sep 17 00:00:00 2001 From: Aykut Bozkurt Date: Tue, 13 Jan 2026 01:37:44 +0300 Subject: [PATCH 2/3] separate structs for parse and analyze Signed-off-by: Aykut Bozkurt --- .../src/data_file/data_file_stats.c | 2 +- .../pg_lake/iceberg/api/partitioning.h | 18 +- .../src/iceberg/partitioning/partition.c | 2 +- .../iceberg/partitioning/spec_generation.c | 2 +- pg_lake_table/src/fdw/data_file_pruning.c | 16 +- pg_lake_table/src/fdw/partition_transform.c | 64 +++---- .../fdw/partitioning/partition_by_parser.c | 158 +++++++++--------- pg_lake_table/src/test/test_partition_tuple.c | 14 +- 8 files changed, 147 insertions(+), 129 deletions(-) diff --git a/pg_lake_engine/src/data_file/data_file_stats.c b/pg_lake_engine/src/data_file/data_file_stats.c index 5e5bbb43a..1d78c82a7 100644 --- a/pg_lake_engine/src/data_file/data_file_stats.c +++ b/pg_lake_engine/src/data_file/data_file_stats.c @@ -256,7 +256,7 @@ ExtractMinMaxForColumn(Datum map, char *colName, List **names, List **mins, List if (minText != NULL && maxText != NULL) { - *names = lappend(*names, pstrdup(colName)); + *names = lappend(*names, colName); *mins = lappend(*mins, minText); *maxs = lappend(*maxs, maxText); } diff --git a/pg_lake_iceberg/include/pg_lake/iceberg/api/partitioning.h b/pg_lake_iceberg/include/pg_lake/iceberg/api/partitioning.h index 535d436c4..6535df6b9 100644 --- a/pg_lake_iceberg/include/pg_lake/iceberg/api/partitioning.h +++ b/pg_lake_iceberg/include/pg_lake/iceberg/api/partitioning.h @@ -39,7 +39,8 @@ typedef enum IcebergPartitionTransformType PARTITION_TRANSFORM_VOID } IcebergPartitionTransformType; -typedef struct IcebergPartitionTransform +/* Represents a parsed partition transform from table's partition_by string option. */ +typedef struct ParsedIcebergPartitionTransform { IcebergPartitionTransformType type; @@ -52,13 +53,22 @@ typedef struct IcebergPartitionTransform size_t truncateLen; }; - IcebergPartitionSpecField *specField; + const char *columnName; +} ParsedIcebergPartitionTransform; + +/* Represents an analyzed partition transform with all necessary info. */ +typedef struct IcebergPartitionTransform +{ + /* parsed transform info */ + ParsedIcebergPartitionTransform parsedTransform; + + /* spec field info */ + IcebergPartitionSpecField specField; /* source field of the column to which transform applies */ - DataFileSchemaField *sourceField; + DataFileSchemaField sourceField; /* Postgres column info to which transform applies */ - const char *columnName; AttrNumber attnum; PGType pgType; diff --git a/pg_lake_iceberg/src/iceberg/partitioning/partition.c b/pg_lake_iceberg/src/iceberg/partitioning/partition.c index 1798e067b..1d054f5b6 100644 --- a/pg_lake_iceberg/src/iceberg/partitioning/partition.c +++ b/pg_lake_iceberg/src/iceberg/partitioning/partition.c @@ -225,7 +225,7 @@ FindPartitionTransformById(List *transforms, int32_t partitionFieldId, bool erro { IcebergPartitionTransform *transform = (IcebergPartitionTransform *) lfirst(cell); - if (transform->specField->field_id == partitionFieldId) + if (transform->specField.field_id == partitionFieldId) return transform; } diff --git a/pg_lake_iceberg/src/iceberg/partitioning/spec_generation.c b/pg_lake_iceberg/src/iceberg/partitioning/spec_generation.c index ccd5762cb..95084b69a 100644 --- a/pg_lake_iceberg/src/iceberg/partitioning/spec_generation.c +++ b/pg_lake_iceberg/src/iceberg/partitioning/spec_generation.c @@ -60,7 +60,7 @@ BuildPartitionSpecFromPartitionTransforms(Oid relationId, List *partitionTransfo { IcebergPartitionTransform *transform = lfirst(transformCell); - spec->fields[fieldIndex] = *(transform->specField); + spec->fields[fieldIndex] = transform->specField; fieldIndex++; } diff --git a/pg_lake_table/src/fdw/data_file_pruning.c b/pg_lake_table/src/fdw/data_file_pruning.c index 10c1b5d3a..aca5d855d 100644 --- a/pg_lake_table/src/fdw/data_file_pruning.c +++ b/pg_lake_table/src/fdw/data_file_pruning.c @@ -676,7 +676,7 @@ GetColumnBoundConstraintsFromPartition(Oid relationId, ColumnToFieldIdMapping * continue; /* skip if transform's sourceId does not match the entry's fieldId */ - if (partitionTransform->sourceField->id != entry->fieldId) + if (partitionTransform->sourceField.id != entry->fieldId) continue; Expr *boundsConstraint = @@ -699,7 +699,7 @@ static Expr * PartitionFieldBoundConstraint(PartitionField * partitionField, IcebergPartitionTransform * partitionTransform, ColumnToFieldIdMapping * entry) { - IcebergPartitionTransformType type = partitionTransform->type; + IcebergPartitionTransformType type = partitionTransform->parsedTransform.type; if (type != PARTITION_TRANSFORM_IDENTITY && partitionField->value == NULL) @@ -751,7 +751,7 @@ IdentityPartitionFieldBoundConstraint(PartitionField * partitionField, { bool isNull = false; Datum partitionDatum = - PartitionValueToDatum(partitionTransform->type, partitionField->value, partitionField->value_length, + PartitionValueToDatum(partitionTransform->parsedTransform.type, partitionField->value, partitionField->value_length, partitionTransform->resultPgType, &isNull); OpExpr *columnBoundEquality = copyObject(entry->equalityOperatorExpression); @@ -780,7 +780,7 @@ TruncatePartitionFieldBoundConstraint(PartitionField * partitionField, if (pgType.postgresTypeOid == INT4OID || pgType.postgresTypeOid == INT2OID) { int32 partitionValue = *(int32_t *) partitionField->value; - int truncateLen = partitionTransform->truncateLen; + int truncateLen = partitionTransform->parsedTransform.truncateLen; int32 upperBound; @@ -798,7 +798,7 @@ TruncatePartitionFieldBoundConstraint(PartitionField * partitionField, else if (pgType.postgresTypeOid == INT8OID) { int64 partitionValue = *(int64_t *) partitionField->value; - int truncateLen = partitionTransform->truncateLen; + int truncateLen = partitionTransform->parsedTransform.truncateLen; int64 upperBound; @@ -825,7 +825,7 @@ TruncatePartitionFieldBoundConstraint(PartitionField * partitionField, return NULL; } - int truncateLen = partitionTransform->truncateLen; + int truncateLen = partitionTransform->parsedTransform.truncateLen; char *truncatedUpperBound = TruncateUpperBoundForText(pstrdup(partitionValue), truncateLen); if (truncatedUpperBound == NULL) @@ -848,7 +848,7 @@ TruncatePartitionFieldBoundConstraint(PartitionField * partitionField, memcpy(VARDATA_ANY(partitionValue), partitionField->value, partitionField->value_length); bytea *partitionValueCopy = (bytea *) pg_detoast_datum_copy((struct varlena *) partitionValue); - int truncateLen = partitionTransform->truncateLen; + int truncateLen = partitionTransform->parsedTransform.truncateLen; /* increment the last byte of the upper bound, which does not overflow */ partitionValueCopy = TruncateUpperBoundForBytea(partitionValueCopy, truncateLen); @@ -1813,7 +1813,7 @@ ExtendClausesForBucketPartitioning(Partition * partition, List *partitionTransfo if (partitionTransform == NULL) continue; - if (partitionTransform->type != PARTITION_TRANSFORM_BUCKET) + if (partitionTransform->parsedTransform.type != PARTITION_TRANSFORM_BUCKET) { /* only extend restrict info for bucket transform */ continue; diff --git a/pg_lake_table/src/fdw/partition_transform.c b/pg_lake_table/src/fdw/partition_transform.c index 4ece3ee4e..f39950858 100644 --- a/pg_lake_table/src/fdw/partition_transform.c +++ b/pg_lake_table/src/fdw/partition_transform.c @@ -152,7 +152,7 @@ PartitionTransformsEqual(IcebergPartitionSpec * spec, List *partitionTransforms) * ErrorIfColumnEverUsedInIcebergPartitionSpec(). Still, let's be * defensive and also check source field ids. */ - if (specField->source_id != transform->sourceField->id) + if (specField->source_id != transform->sourceField.id) return false; /* @@ -162,7 +162,7 @@ PartitionTransformsEqual(IcebergPartitionSpec * spec, List *partitionTransforms) * Iceberg does here: * https://github.com/apache/iceberg/blob/8b55ac834015ce664f879ecfe1e80a941a994420/api/src/main/java/org/apache/iceberg/PartitionSpec.java#L239-L259 */ - if (strcasecmp(specField->name, transform->specField->name) != 0) + if (strcasecmp(specField->name, transform->specField.name) != 0) { return false; } @@ -251,16 +251,18 @@ GetPartitionTransformFromSpecField(Oid relationId, IcebergPartitionSpecField * s { IcebergPartitionTransform *transform = palloc0(sizeof(IcebergPartitionTransform)); - transform->specField = specField; + transform->specField = *specField; transform->attnum = GetAttributeForFieldId(relationId, specField->source_id); - transform->columnName = get_attname(relationId, transform->attnum, false); + transform->parsedTransform.columnName = get_attname(relationId, transform->attnum, false); transform->pgType = GetAttributePGType(relationId, transform->attnum); if (IsInternalIcebergTable(relationId)) { - transform->sourceField = GetRegisteredFieldForAttribute(relationId, transform->attnum); + DataFileSchemaField *sourceField = GetRegisteredFieldForAttribute(relationId, transform->attnum); + + transform->sourceField = *sourceField; } else { @@ -268,14 +270,16 @@ GetPartitionTransformFromSpecField(Oid relationId, IcebergPartitionSpecField * s DataFileSchema *schema = GetDataFileSchemaForTable(relationId); - transform->sourceField = GetDataFileSchemaFieldById(schema, specField->source_id); + DataFileSchemaField *sourceField = GetDataFileSchemaFieldById(schema, specField->source_id); + + transform->sourceField = *sourceField; } /* parse transform name */ - ParseTransformName(transform->specField->transform, - &transform->type, - &transform->bucketCount, - &transform->truncateLen); + ParseTransformName(transform->specField.transform, + &transform->parsedTransform.type, + &transform->parsedTransform.bucketCount, + &transform->parsedTransform.truncateLen); /* set transform's postgres type */ transform->resultPgType = GetTransformResultPGType(transform); @@ -411,13 +415,13 @@ ApplyPartitionTransformToTuple(IcebergPartitionTransform * transform, TupleTable { PartitionField *field = palloc0(sizeof(PartitionField)); - field->field_name = pstrdup(transform->specField->name); - field->field_id = transform->specField->field_id; + field->field_name = pstrdup(transform->specField.name); + field->field_id = transform->specField.field_id; bool isNull = false; Datum columnValue = slot_getattr(slot, transform->attnum, &isNull); - switch (transform->type) + switch (transform->parsedTransform.type) { case PARTITION_TRANSFORM_IDENTITY: field->value = ApplyIdentityTransformToColumn(transform, columnValue, isNull, @@ -451,7 +455,7 @@ ApplyPartitionTransformToTuple(IcebergPartitionTransform * transform, TupleTable ereport(ERROR, (errcode(ERRCODE_FEATURE_NOT_SUPPORTED), errmsg("applying transform %s is not yet support ", - transform->specField->transform))); + transform->specField.transform))); } field->value_type = GetTransformResultAvroType(transform); @@ -475,7 +479,7 @@ ApplyIdentityTransformToColumn(IcebergPartitionTransform * transform, Datum colu return NULL; } - return PGIcebergBinarySerializePartitionFieldValue(columnValue, transform->sourceField->type, + return PGIcebergBinarySerializePartitionFieldValue(columnValue, transform->sourceField.type, transform->pgType, valueSize); } @@ -496,7 +500,7 @@ ApplyTruncateTransformToColumn(IcebergPartitionTransform * transform, Datum colu PGType sourceType = transform->pgType; PGType resultType = transform->resultPgType; - int64_t truncateLen = (int64_t) transform->truncateLen; + int64_t truncateLen = (int64_t) transform->parsedTransform.truncateLen; Datum truncatedColumnValue = 0; if (sourceType.postgresTypeOid == INT2OID) @@ -765,7 +769,7 @@ ApplyBucketTransformToColumn(IcebergPartitionTransform * transform, Datum column { int64_t value = (int64_t) DatumGetInt16(columnValue); - *bucketValue = (MurmurHash3_32_Long(value) & INT32_MAX) % transform->bucketCount; + *bucketValue = (MurmurHash3_32_Long(value) & INT32_MAX) % transform->parsedTransform.bucketCount; } else if (transform->pgType.postgresTypeOid == INT4OID) { @@ -775,13 +779,13 @@ ApplyBucketTransformToColumn(IcebergPartitionTransform * transform, Datum column */ int64_t value = (int64_t) DatumGetInt32(columnValue); - *bucketValue = (MurmurHash3_32_Long(value) & INT32_MAX) % transform->bucketCount; + *bucketValue = (MurmurHash3_32_Long(value) & INT32_MAX) % transform->parsedTransform.bucketCount; } else if (transform->pgType.postgresTypeOid == INT8OID) { int64_t value = DatumGetInt64(columnValue); - *bucketValue = (MurmurHash3_32_Long(value) & INT32_MAX) % transform->bucketCount; + *bucketValue = (MurmurHash3_32_Long(value) & INT32_MAX) % transform->parsedTransform.bucketCount; } else if (transform->pgType.postgresTypeOid == TEXTOID || transform->pgType.postgresTypeOid == VARCHAROID || @@ -789,13 +793,13 @@ ApplyBucketTransformToColumn(IcebergPartitionTransform * transform, Datum column { const char *value = TextDatumGetCString(columnValue); - *bucketValue = (MurmurHash3_32_Bytes(value, strlen(value)) & INT32_MAX) % transform->bucketCount; + *bucketValue = (MurmurHash3_32_Bytes(value, strlen(value)) & INT32_MAX) % transform->parsedTransform.bucketCount; } else if (transform->pgType.postgresTypeOid == BYTEAOID) { bytea *value = DatumGetByteaP(columnValue); - *bucketValue = (MurmurHash3_32_Bytes(VARDATA_ANY(value), VARSIZE_ANY_EXHDR(value)) & INT32_MAX) % transform->bucketCount; + *bucketValue = (MurmurHash3_32_Bytes(VARDATA_ANY(value), VARSIZE_ANY_EXHDR(value)) & INT32_MAX) % transform->parsedTransform.bucketCount; } else if (transform->pgType.postgresTypeOid == DATEOID) { @@ -807,7 +811,7 @@ ApplyBucketTransformToColumn(IcebergPartitionTransform * transform, Datum column * spec normally hashes int bytes for date type but spark hashes long * bytes of date. We follow spark here. */ - *bucketValue = (MurmurHash3_32_Long(daysFromEpoch) & INT32_MAX) % transform->bucketCount; + *bucketValue = (MurmurHash3_32_Long(daysFromEpoch) & INT32_MAX) % transform->parsedTransform.bucketCount; } else if (transform->pgType.postgresTypeOid == TIMESTAMPOID) { @@ -815,7 +819,7 @@ ApplyBucketTransformToColumn(IcebergPartitionTransform * transform, Datum column int64_t microsecsFromEpoch = AdjustTimestampFromPostgresToUnix(value); - *bucketValue = (MurmurHash3_32_Long(microsecsFromEpoch) & INT32_MAX) % transform->bucketCount; + *bucketValue = (MurmurHash3_32_Long(microsecsFromEpoch) & INT32_MAX) % transform->parsedTransform.bucketCount; } else if (transform->pgType.postgresTypeOid == TIMESTAMPTZOID) { @@ -823,7 +827,7 @@ ApplyBucketTransformToColumn(IcebergPartitionTransform * transform, Datum column int64_t microsecsFromEpoch = AdjustTimestampFromPostgresToUnix(value); - *bucketValue = (MurmurHash3_32_Long(microsecsFromEpoch) & INT32_MAX) % transform->bucketCount; + *bucketValue = (MurmurHash3_32_Long(microsecsFromEpoch) & INT32_MAX) % transform->parsedTransform.bucketCount; } else if (transform->pgType.postgresTypeOid == TIMEOID) { @@ -831,23 +835,23 @@ ApplyBucketTransformToColumn(IcebergPartitionTransform * transform, Datum column int64_t microsecsFromMidnight = value; - *bucketValue = (MurmurHash3_32_Long(microsecsFromMidnight) & INT32_MAX) % transform->bucketCount; + *bucketValue = (MurmurHash3_32_Long(microsecsFromMidnight) & INT32_MAX) % transform->parsedTransform.bucketCount; } else if (transform->pgType.postgresTypeOid == UUIDOID) { size_t valueSize = 0; - unsigned char *value = PGIcebergBinarySerializePartitionFieldValue(columnValue, transform->sourceField->type, + unsigned char *value = PGIcebergBinarySerializePartitionFieldValue(columnValue, transform->sourceField.type, transform->pgType, &valueSize); - *bucketValue = (MurmurHash3_32_Bytes(value, valueSize) & INT32_MAX) % transform->bucketCount; + *bucketValue = (MurmurHash3_32_Bytes(value, valueSize) & INT32_MAX) % transform->parsedTransform.bucketCount; } else if (transform->pgType.postgresTypeOid == NUMERICOID) { size_t valueSize = 0; - unsigned char *value = PGIcebergBinarySerializePartitionFieldValue(columnValue, transform->sourceField->type, + unsigned char *value = PGIcebergBinarySerializePartitionFieldValue(columnValue, transform->sourceField.type, transform->pgType, &valueSize); - *bucketValue = (MurmurHash3_32_Bytes(value, valueSize) & INT32_MAX) % transform->bucketCount; + *bucketValue = (MurmurHash3_32_Bytes(value, valueSize) & INT32_MAX) % transform->parsedTransform.bucketCount; } else { @@ -975,7 +979,7 @@ SerializePartitionValueToPGText(void *value, size_t valueLength, IcebergPartitio /* First, deserialize back */ bool isNull = false; Datum partitionDatum = - PartitionValueToDatum(transform->type, value, valueLength, + PartitionValueToDatum(transform->parsedTransform.type, value, valueLength, transform->resultPgType, &isNull); if (isNull) diff --git a/pg_lake_table/src/fdw/partitioning/partition_by_parser.c b/pg_lake_table/src/fdw/partitioning/partition_by_parser.c index 697c487dc..98be72835 100644 --- a/pg_lake_table/src/fdw/partitioning/partition_by_parser.c +++ b/pg_lake_table/src/fdw/partitioning/partition_by_parser.c @@ -43,12 +43,12 @@ static void SkipWhitespace(const char **ptr); static int ParseInteger(const char **ptr); static char *ParseIdentifier(const char **ptr); static IcebergPartitionTransformType LookupTransformType(const char *str); -static IcebergPartitionTransform * ParseOneTransform(const char **ptr); +static ParsedIcebergPartitionTransform * ParseOneTransform(const char **ptr); static const char *NormalizeTransformColumnName(const char *colName); /* analyzer helpers */ -static void EnsureTransformSourceColumnExists(IcebergPartitionTransform * transform, Oid relationId); -static void EnsureTransformSourceColumnScalar(IcebergPartitionTransform * transform, DataFileSchemaField * sourceField); +static void EnsureTransformSourceColumnExists(ParsedIcebergPartitionTransform * parsedTransform, Oid relationId); +static void EnsureTransformSourceColumnScalar(ParsedIcebergPartitionTransform * parsedTransform, DataFileSchemaField * sourceField); static void EnsureNoDuplicateTransforms(List *transforms); static void EnsureValidTypeForTransform(IcebergPartitionTransformType transformType, Oid typeOid); static void EnsureValidTypeForIdentityTransform(Oid typeOid); @@ -84,7 +84,7 @@ ParseIcebergTablePartitionBy(Oid relationId) Assert(partitionBy != NULL); - List *transforms = NIL; + List *parsedTransforms = NIL; bool seenComma = false; @@ -109,9 +109,9 @@ ParseIcebergTablePartitionBy(Oid relationId) } /* Parse one transform */ - IcebergPartitionTransform *transform = ParseOneTransform(&p); + ParsedIcebergPartitionTransform *parsedTransform = ParseOneTransform(&p); - transforms = lappend(transforms, transform); + parsedTransforms = lappend(parsedTransforms, parsedTransform); /* after one transform, expect either comma or end of string */ SkipWhitespace(&p); @@ -138,7 +138,7 @@ ParseIcebergTablePartitionBy(Oid relationId) } } - return transforms; + return parsedTransforms; } @@ -159,10 +159,10 @@ IsIcebergTableWithDefaultPartitionSpec(Oid relationId) * Parse one transform from the current position of the string. * (ex: "year(col)", "bucket(mycol, 10)", or just "colName") */ -static IcebergPartitionTransform * +static ParsedIcebergPartitionTransform * ParseOneTransform(const char **ptr) { - IcebergPartitionTransform *transform = palloc0(sizeof(IcebergPartitionTransform)); + ParsedIcebergPartitionTransform *parsedTransform = palloc0(sizeof(ParsedIcebergPartitionTransform)); /* Step 1: parse a token => either transform name or column name. */ char *token = ParseIdentifier(ptr); @@ -181,13 +181,13 @@ ParseOneTransform(const char **ptr) if (**ptr == '(') { /* So 'token' is transform name (year, month, etc.) */ - transform->type = LookupTransformType(token); + parsedTransform->type = LookupTransformType(token); (*ptr)++; /* skip '(' */ SkipWhitespace(ptr); - if (transform->type == PARTITION_TRANSFORM_BUCKET || - transform->type == PARTITION_TRANSFORM_TRUNCATE) + if (parsedTransform->type == PARTITION_TRANSFORM_BUCKET || + parsedTransform->type == PARTITION_TRANSFORM_TRUNCATE) { int param = ParseInteger(ptr); @@ -198,17 +198,17 @@ ParseOneTransform(const char **ptr) { (*ptr)++; - if (transform->type == PARTITION_TRANSFORM_BUCKET) - transform->bucketCount = (size_t) param; + if (parsedTransform->type == PARTITION_TRANSFORM_BUCKET) + parsedTransform->bucketCount = (size_t) param; else - transform->truncateLen = (size_t) param; + parsedTransform->truncateLen = (size_t) param; } else { ereport(ERROR, (errcode(ERRCODE_SYNTAX_ERROR), errmsg("expected comma after column name for \"%s\" transform", - (transform->type == PARTITION_TRANSFORM_BUCKET) ? "bucket" : "truncate"))); + (parsedTransform->type == PARTITION_TRANSFORM_BUCKET) ? "bucket" : "truncate"))); } } @@ -217,7 +217,7 @@ ParseOneTransform(const char **ptr) /* Next parse the column name inside parentheses. */ char *colName = ParseIdentifier(ptr); - transform->columnName = NormalizeTransformColumnName(colName); + parsedTransform->columnName = NormalizeTransformColumnName(colName); SkipWhitespace(ptr); if (**ptr == ')') @@ -234,11 +234,11 @@ ParseOneTransform(const char **ptr) else { /* The token is the column name; interpret as IDENTITY transform. */ - transform->type = PARTITION_TRANSFORM_IDENTITY; - transform->columnName = NormalizeTransformColumnName(token); + parsedTransform->type = PARTITION_TRANSFORM_IDENTITY; + parsedTransform->columnName = NormalizeTransformColumnName(token); } - return transform; + return parsedTransform; } @@ -480,32 +480,38 @@ NormalizeTransformColumnName(const char *colName) * and also checks for duplicate transforms on the same source id in the partition spec. */ List * -AnalyzeIcebergTablePartitionBy(Oid relationId, List *transforms) +AnalyzeIcebergTablePartitionBy(Oid relationId, List *parsedTransforms) { int largestPartitionFieldId = GetLargestPartitionFieldId(relationId); + List *analyzedTransforms = NIL; + /* analyze */ - ListCell *transformCell = NULL; + ListCell *parsedTransformCell = NULL; - foreach(transformCell, transforms) + foreach(parsedTransformCell, parsedTransforms) { - IcebergPartitionTransform *transform = lfirst(transformCell); + ParsedIcebergPartitionTransform *parsedTransform = lfirst(parsedTransformCell); /* * 1) Look up the column in the relation's TupleDesc. We do a * case-sensitive or case-insensitive match depending on your FDW's * requirements. Here we do an exact match. */ - EnsureTransformSourceColumnExists(transform, relationId); + EnsureTransformSourceColumnExists(parsedTransform, relationId); + + IcebergPartitionTransform *analyzedTransform = palloc0(sizeof(IcebergPartitionTransform)); + + analyzedTransform->parsedTransform = *parsedTransform; /* set column no */ - transform->attnum = get_attnum(relationId, transform->columnName); + analyzedTransform->attnum = get_attnum(relationId, parsedTransform->columnName); Oid collation = InvalidOid; - get_atttypetypmodcoll(relationId, transform->attnum, - &transform->pgType.postgresTypeOid, - &transform->pgType.postgresTypeMod, + get_atttypetypmodcoll(relationId, analyzedTransform->attnum, + &analyzedTransform->pgType.postgresTypeOid, + &analyzedTransform->pgType.postgresTypeMod, &collation); /* set column type */ @@ -515,49 +521,47 @@ AnalyzeIcebergTablePartitionBy(Oid relationId, List *transforms) ereport(ERROR, (errcode(ERRCODE_FEATURE_NOT_SUPPORTED), errmsg("Columns with collation are not supported in partition_by: \"%s\"", - transform->columnName))); + parsedTransform->columnName))); } /* set result column type */ - transform->resultPgType = GetTransformResultPGType(transform); - + analyzedTransform->resultPgType = GetTransformResultPGType(analyzedTransform); /* set source field id */ - DataFileSchemaField *sourceField = GetRegisteredFieldForAttribute(relationId, transform->attnum); + DataFileSchemaField *sourceField = GetRegisteredFieldForAttribute(relationId, analyzedTransform->attnum); /* 2) Check scalar column */ - EnsureTransformSourceColumnScalar(transform, sourceField); - - transform->sourceField = sourceField; + EnsureTransformSourceColumnScalar(parsedTransform, sourceField); - transform->specField = palloc0(sizeof(IcebergPartitionSpecField)); + analyzedTransform->sourceField = *sourceField; /* set transform name */ - transform->specField->transform = GenerateTransformName(transform); - transform->specField->transform_length = strlen(transform->specField->transform); + analyzedTransform->specField.transform = GenerateTransformName(analyzedTransform); + analyzedTransform->specField.transform_length = strlen(analyzedTransform->specField.transform); /* set partition field name */ - transform->specField->name = GeneratePartitionFieldName(transform, relationId); - transform->specField->name_length = strlen(transform->specField->name); + analyzedTransform->specField.name = GeneratePartitionFieldName(analyzedTransform, relationId); + analyzedTransform->specField.name_length = strlen(analyzedTransform->specField.name); /* set partition field id */ - transform->specField->field_id = ++largestPartitionFieldId; - + analyzedTransform->specField.field_id = ++largestPartitionFieldId; /* set source field id */ - transform->specField->source_id = sourceField->id; - transform->specField->source_ids_length = 1; - transform->specField->source_ids = palloc0(sizeof(int) * transform->specField->source_ids_length); - transform->specField->source_ids[0] = transform->specField->source_id; + analyzedTransform->specField.source_id = sourceField->id; + analyzedTransform->specField.source_ids_length = 1; + analyzedTransform->specField.source_ids = palloc0(sizeof(int) * analyzedTransform->specField.source_ids_length); + analyzedTransform->specField.source_ids[0] = analyzedTransform->specField.source_id; /* 3) Check column type compatibility. */ - EnsureValidTypeForTransform(transform->type, transform->pgType.postgresTypeOid); + EnsureValidTypeForTransform(analyzedTransform->parsedTransform.type, analyzedTransform->pgType.postgresTypeOid); + + analyzedTransforms = lappend(analyzedTransforms, analyzedTransform); } /* - * 4) Check for duplicate transoforms on the same source id. + * 4) Check for duplicate transforms on the same source id. */ - EnsureNoDuplicateTransforms(transforms); + EnsureNoDuplicateTransforms(analyzedTransforms); - return transforms; + return analyzedTransforms; } @@ -599,14 +603,14 @@ GetIcebergTablePartitionByOption(Oid relationId) static const char * GenerateTransformName(IcebergPartitionTransform * transform) { - switch (transform->type) + switch (transform->parsedTransform.type) { case PARTITION_TRANSFORM_IDENTITY: return "identity"; case PARTITION_TRANSFORM_BUCKET: - return psprintf("bucket[%zu]", transform->bucketCount); + return psprintf("bucket[%zu]", transform->parsedTransform.bucketCount); case PARTITION_TRANSFORM_TRUNCATE: - return psprintf("truncate[%zu]", transform->truncateLen); + return psprintf("truncate[%zu]", transform->parsedTransform.truncateLen); case PARTITION_TRANSFORM_YEAR: return "year"; case PARTITION_TRANSFORM_MONTH: @@ -620,7 +624,7 @@ GenerateTransformName(IcebergPartitionTransform * transform) default: ereport(ERROR, (errcode(ERRCODE_FEATURE_NOT_SUPPORTED), - errmsg("unknown partition transform type %d", transform->type))); + errmsg("unknown partition transform type %d", transform->parsedTransform.type))); } } @@ -631,28 +635,28 @@ GenerateTransformName(IcebergPartitionTransform * transform) static const char * GeneratePartitionFieldName(IcebergPartitionTransform * transform, Oid relationId) { - switch (transform->type) + switch (transform->parsedTransform.type) { case PARTITION_TRANSFORM_IDENTITY: - return transform->columnName; + return transform->parsedTransform.columnName; case PARTITION_TRANSFORM_BUCKET: - return psprintf("%s_bucket_%zu", transform->columnName, transform->bucketCount); + return psprintf("%s_bucket_%zu", transform->parsedTransform.columnName, transform->parsedTransform.bucketCount); case PARTITION_TRANSFORM_TRUNCATE: - return psprintf("%s_trunc_%zu", transform->columnName, transform->truncateLen); + return psprintf("%s_trunc_%zu", transform->parsedTransform.columnName, transform->parsedTransform.truncateLen); case PARTITION_TRANSFORM_YEAR: - return psprintf("%s_year", transform->columnName); + return psprintf("%s_year", transform->parsedTransform.columnName); case PARTITION_TRANSFORM_MONTH: - return psprintf("%s_month", transform->columnName); + return psprintf("%s_month", transform->parsedTransform.columnName); case PARTITION_TRANSFORM_DAY: - return psprintf("%s_day", transform->columnName); + return psprintf("%s_day", transform->parsedTransform.columnName); case PARTITION_TRANSFORM_HOUR: - return psprintf("%s_hour", transform->columnName); + return psprintf("%s_hour", transform->parsedTransform.columnName); case PARTITION_TRANSFORM_VOID: - return psprintf("%s_void", transform->columnName); + return psprintf("%s_void", transform->parsedTransform.columnName); default: ereport(ERROR, (errcode(ERRCODE_FEATURE_NOT_SUPPORTED), - errmsg("unknown partition transform type %d", transform->type))); + errmsg("unknown partition transform type %d", transform->parsedTransform.type))); } } @@ -662,7 +666,7 @@ GeneratePartitionFieldName(IcebergPartitionTransform * transform, Oid relationId * given transform exists in the relation. If it doesn't, it raises an error. */ static void -EnsureTransformSourceColumnExists(IcebergPartitionTransform * transform, Oid relationId) +EnsureTransformSourceColumnExists(ParsedIcebergPartitionTransform * parsedTransform, Oid relationId) { Relation rel; TupleDesc desc; @@ -684,7 +688,7 @@ EnsureTransformSourceColumnExists(IcebergPartitionTransform * transform, Oid rel if (attr->attisdropped) continue; /* skip dropped columns */ - if (strcmp(NameStr(attr->attname), transform->columnName) == 0) + if (strcmp(NameStr(attr->attname), parsedTransform->columnName) == 0) { foundCol = true; break; @@ -697,7 +701,7 @@ EnsureTransformSourceColumnExists(IcebergPartitionTransform * transform, Oid rel ereport(ERROR, (errcode(ERRCODE_UNDEFINED_COLUMN), errmsg("column \"%s\" does not exist in relation \"%s\"", - transform->columnName, + parsedTransform->columnName, get_rel_name(relationId)))); } @@ -714,7 +718,7 @@ GetTransformResultPGType(IcebergPartitionTransform * transform) Oid resultType = InvalidOid; int32_t resultTypMod = -1; - switch (transform->type) + switch (transform->parsedTransform.type) { case PARTITION_TRANSFORM_IDENTITY: resultType = transform->pgType.postgresTypeOid; @@ -752,7 +756,7 @@ GetTransformResultPGType(IcebergPartitionTransform * transform) default: ereport(ERROR, (errcode(ERRCODE_FEATURE_NOT_SUPPORTED), - errmsg("unknown partition transform type %d", transform->type))); + errmsg("unknown partition transform type %d", transform->parsedTransform.type))); break; } @@ -783,10 +787,10 @@ AdjustTypmodForTruncateTransformIfNeeded(IcebergPartitionTransform * transform) * - bpchar(N): any smaller (< N) string will be space-padded to N. */ if ((pgType.postgresTypeOid == VARCHAROID || pgType.postgresTypeOid == BPCHAROID) && - pgType.postgresTypeMod != -1 && GetAnyCharLengthFrom(pgType.postgresTypeMod) > transform->truncateLen) + pgType.postgresTypeMod != -1 && GetAnyCharLengthFrom(pgType.postgresTypeMod) > transform->parsedTransform.truncateLen) { resultTypMod = AdjustAnyCharTypmod(transform->pgType.postgresTypeMod, - transform->truncateLen); + transform->parsedTransform.truncateLen); } else { @@ -802,13 +806,13 @@ AdjustTypmodForTruncateTransformIfNeeded(IcebergPartitionTransform * transform) * given transform is a scalar type. If it isn't, it raises an error. */ static void -EnsureTransformSourceColumnScalar(IcebergPartitionTransform * transform, DataFileSchemaField * sourceField) +EnsureTransformSourceColumnScalar(ParsedIcebergPartitionTransform * parsedTransform, DataFileSchemaField * sourceField) { if (sourceField->type->type != FIELD_TYPE_SCALAR) ereport(ERROR, (errcode(ERRCODE_DATATYPE_MISMATCH), errmsg("partition transform column \"%s\" must be a scalar type", - transform->columnName))); + parsedTransform->columnName))); } @@ -827,12 +831,12 @@ EnsureNoDuplicateTransforms(List *transforms) { IcebergPartitionTransform *other = list_nth(transforms, j); - if (transform->type == other->type && - transform->sourceField->id == other->sourceField->id) + if (transform->parsedTransform.type == other->parsedTransform.type && + transform->sourceField.id == other->sourceField.id) ereport(ERROR, (errcode(ERRCODE_DUPLICATE_OBJECT), errmsg("\"%s\" transform on column \"%s\" appears multiple times in partition spec", - transform->specField->transform, transform->columnName))); + transform->specField.transform, transform->parsedTransform.columnName))); } } } diff --git a/pg_lake_table/src/test/test_partition_tuple.c b/pg_lake_table/src/test/test_partition_tuple.c index 673979b6c..b4b3cc56a 100644 --- a/pg_lake_table/src/test/test_partition_tuple.c +++ b/pg_lake_table/src/test/test_partition_tuple.c @@ -97,19 +97,19 @@ get_partition_tuple(PG_FUNCTION_ARGS) /* corresponding transform for partition field */ IcebergPartitionTransform *transform = list_nth(transforms, i); - AttrNumber attnum = get_attnum(relationId, transform->columnName); + AttrNumber attnum = get_attnum(relationId, transform->parsedTransform.columnName); if (attnum == InvalidAttrNumber) { ereport(ERROR, (errcode(ERRCODE_UNDEFINED_COLUMN), errmsg("column %s does not exist in relation %u", - transform->columnName, relationId))); + transform->parsedTransform.columnName, relationId))); } PGType pgType = transform->resultPgType; - values[i] = PartitionValueToDatum(transform->type, field->value, field->value_length, pgType, &nulls[i]); + values[i] = PartitionValueToDatum(transform->parsedTransform.type, field->value, field->value_length, pgType, &nulls[i]); } ExecDropSingleTupleTableSlot(slot); @@ -201,7 +201,7 @@ get_partition_summary(PG_FUNCTION_ARGS) partitionSummary->lower_bound_length, transform); - values[2] = AdjustFieldSummaryTextToSpark(lowerBoundText, sourceType, transform->type); + values[2] = AdjustFieldSummaryTextToSpark(lowerBoundText, sourceType, transform->parsedTransform.type); nulls[2] = false; } else @@ -215,7 +215,7 @@ get_partition_summary(PG_FUNCTION_ARGS) partitionSummary->upper_bound_length, transform); - values[3] = AdjustFieldSummaryTextToSpark(upperBoundText, sourceType, transform->type); + values[3] = AdjustFieldSummaryTextToSpark(upperBoundText, sourceType, transform->parsedTransform.type); nulls[3] = false; } else @@ -296,12 +296,12 @@ EnsureTupleDescMatchTransforms(TupleDesc tupledesc, List *transforms) (errcode(ERRCODE_INVALID_PARAMETER_VALUE), errmsg("attribute %s's type %s in tupledesc does not match transform %s's type %s", attname, format_type_be(column->atttypid), - transform->columnName, format_type_be(transformResultPgType.postgresTypeOid)))); + transform->parsedTransform.columnName, format_type_be(transformResultPgType.postgresTypeOid)))); if (column->atttypmod != transformResultPgType.postgresTypeMod) ereport(ERROR, (errcode(ERRCODE_INVALID_PARAMETER_VALUE), errmsg("attribute %s's typmod %d in tupledesc does not match transform %s's typmod %d", - attname, column->atttypmod, transform->columnName, + attname, column->atttypmod, transform->parsedTransform.columnName, transformResultPgType.postgresTypeMod))); } } From d7db9918ff6e4a9130367bd5e647f628dd1b9c90 Mon Sep 17 00:00:00 2001 From: Aykut Bozkurt Date: Mon, 26 Jan 2026 15:33:08 +0300 Subject: [PATCH 3/3] deep copy spec field Signed-off-by: Aykut Bozkurt --- .../pg_lake/iceberg/api/partitioning.h | 4 +-- .../include/pg_lake/iceberg/metadata_spec.h | 1 + .../src/iceberg/partitioning/partition.c | 2 +- .../iceberg/partitioning/spec_generation.c | 2 +- .../src/iceberg/write_table_metadata.c | 25 +++++++++++++++++ pg_lake_table/src/fdw/data_file_pruning.c | 2 +- pg_lake_table/src/fdw/partition_transform.c | 28 ++++++++----------- .../fdw/partitioning/partition_by_parser.c | 27 +++++++++--------- 8 files changed, 57 insertions(+), 34 deletions(-) diff --git a/pg_lake_iceberg/include/pg_lake/iceberg/api/partitioning.h b/pg_lake_iceberg/include/pg_lake/iceberg/api/partitioning.h index 6535df6b9..dff8e66fe 100644 --- a/pg_lake_iceberg/include/pg_lake/iceberg/api/partitioning.h +++ b/pg_lake_iceberg/include/pg_lake/iceberg/api/partitioning.h @@ -63,10 +63,10 @@ typedef struct IcebergPartitionTransform ParsedIcebergPartitionTransform parsedTransform; /* spec field info */ - IcebergPartitionSpecField specField; + IcebergPartitionSpecField *specField; /* source field of the column to which transform applies */ - DataFileSchemaField sourceField; + DataFileSchemaField *sourceField; /* Postgres column info to which transform applies */ AttrNumber attnum; diff --git a/pg_lake_iceberg/include/pg_lake/iceberg/metadata_spec.h b/pg_lake_iceberg/include/pg_lake/iceberg/metadata_spec.h index 03afdc53c..9f31d2ebc 100644 --- a/pg_lake_iceberg/include/pg_lake/iceberg/metadata_spec.h +++ b/pg_lake_iceberg/include/pg_lake/iceberg/metadata_spec.h @@ -292,3 +292,4 @@ extern PGDLLEXPORT IcebergTableMetadata * ReadIcebergTableMetadata(const char *t extern PGDLLEXPORT char *WriteIcebergTableMetadataToJson(IcebergTableMetadata * metadata); extern PGDLLEXPORT void AppendIcebergTableSchemaForRestCatalog(StringInfo command, IcebergTableSchema * schemas, size_t schemas_length); extern PGDLLEXPORT void AppendIcebergPartitionSpecFields(StringInfo command, IcebergPartitionSpecField * fields, size_t fields_length); +extern PGDLLEXPORT IcebergPartitionSpecField * DeepCopyIcebergPartitionSpecField(const IcebergPartitionSpecField * field); diff --git a/pg_lake_iceberg/src/iceberg/partitioning/partition.c b/pg_lake_iceberg/src/iceberg/partitioning/partition.c index 1d054f5b6..1798e067b 100644 --- a/pg_lake_iceberg/src/iceberg/partitioning/partition.c +++ b/pg_lake_iceberg/src/iceberg/partitioning/partition.c @@ -225,7 +225,7 @@ FindPartitionTransformById(List *transforms, int32_t partitionFieldId, bool erro { IcebergPartitionTransform *transform = (IcebergPartitionTransform *) lfirst(cell); - if (transform->specField.field_id == partitionFieldId) + if (transform->specField->field_id == partitionFieldId) return transform; } diff --git a/pg_lake_iceberg/src/iceberg/partitioning/spec_generation.c b/pg_lake_iceberg/src/iceberg/partitioning/spec_generation.c index 95084b69a..90c50b286 100644 --- a/pg_lake_iceberg/src/iceberg/partitioning/spec_generation.c +++ b/pg_lake_iceberg/src/iceberg/partitioning/spec_generation.c @@ -60,7 +60,7 @@ BuildPartitionSpecFromPartitionTransforms(Oid relationId, List *partitionTransfo { IcebergPartitionTransform *transform = lfirst(transformCell); - spec->fields[fieldIndex] = transform->specField; + spec->fields[fieldIndex] = *(DeepCopyIcebergPartitionSpecField(transform->specField)); fieldIndex++; } diff --git a/pg_lake_iceberg/src/iceberg/write_table_metadata.c b/pg_lake_iceberg/src/iceberg/write_table_metadata.c index 0c1650cbb..8a3e2c17c 100644 --- a/pg_lake_iceberg/src/iceberg/write_table_metadata.c +++ b/pg_lake_iceberg/src/iceberg/write_table_metadata.c @@ -764,6 +764,31 @@ AppendIcebergPartitionSpecFields(StringInfo command, IcebergPartitionSpecField * appendStringInfoString(command, "]"); } +/* + * DeepCopyIcebergPartitionSpecField deep copies a IcebergPartitionSpecField. + */ +IcebergPartitionSpecField * +DeepCopyIcebergPartitionSpecField(const IcebergPartitionSpecField * field) +{ + IcebergPartitionSpecField *copiedField = palloc0(sizeof(IcebergPartitionSpecField)); + + copiedField->field_id = field->field_id; + copiedField->name = pstrdup(field->name); + copiedField->name_length = field->name_length; + copiedField->transform = pstrdup(field->transform); + copiedField->transform_length = field->transform_length; + + copiedField->source_id = field->source_id; + copiedField->source_ids_length = field->source_ids_length; + if (field->source_ids_length > 0) + { + copiedField->source_ids = palloc0(field->source_ids_length * sizeof(int32)); + memcpy(copiedField->source_ids, field->source_ids, field->source_ids_length * sizeof(int32)); + } + + return copiedField; +} + static void AppendIcebergSortOrderFields(StringInfo command, IcebergSortOrderField * fields, size_t fields_length) { diff --git a/pg_lake_table/src/fdw/data_file_pruning.c b/pg_lake_table/src/fdw/data_file_pruning.c index aca5d855d..77a88a558 100644 --- a/pg_lake_table/src/fdw/data_file_pruning.c +++ b/pg_lake_table/src/fdw/data_file_pruning.c @@ -676,7 +676,7 @@ GetColumnBoundConstraintsFromPartition(Oid relationId, ColumnToFieldIdMapping * continue; /* skip if transform's sourceId does not match the entry's fieldId */ - if (partitionTransform->sourceField.id != entry->fieldId) + if (partitionTransform->sourceField->id != entry->fieldId) continue; Expr *boundsConstraint = diff --git a/pg_lake_table/src/fdw/partition_transform.c b/pg_lake_table/src/fdw/partition_transform.c index f39950858..2e650a4f7 100644 --- a/pg_lake_table/src/fdw/partition_transform.c +++ b/pg_lake_table/src/fdw/partition_transform.c @@ -152,7 +152,7 @@ PartitionTransformsEqual(IcebergPartitionSpec * spec, List *partitionTransforms) * ErrorIfColumnEverUsedInIcebergPartitionSpec(). Still, let's be * defensive and also check source field ids. */ - if (specField->source_id != transform->sourceField.id) + if (specField->source_id != transform->sourceField->id) return false; /* @@ -162,7 +162,7 @@ PartitionTransformsEqual(IcebergPartitionSpec * spec, List *partitionTransforms) * Iceberg does here: * https://github.com/apache/iceberg/blob/8b55ac834015ce664f879ecfe1e80a941a994420/api/src/main/java/org/apache/iceberg/PartitionSpec.java#L239-L259 */ - if (strcasecmp(specField->name, transform->specField.name) != 0) + if (strcasecmp(specField->name, transform->specField->name) != 0) { return false; } @@ -251,7 +251,7 @@ GetPartitionTransformFromSpecField(Oid relationId, IcebergPartitionSpecField * s { IcebergPartitionTransform *transform = palloc0(sizeof(IcebergPartitionTransform)); - transform->specField = *specField; + transform->specField = DeepCopyIcebergPartitionSpecField(specField); transform->attnum = GetAttributeForFieldId(relationId, specField->source_id); @@ -260,9 +260,7 @@ GetPartitionTransformFromSpecField(Oid relationId, IcebergPartitionSpecField * s if (IsInternalIcebergTable(relationId)) { - DataFileSchemaField *sourceField = GetRegisteredFieldForAttribute(relationId, transform->attnum); - - transform->sourceField = *sourceField; + transform->sourceField = GetRegisteredFieldForAttribute(relationId, transform->attnum); } else { @@ -270,13 +268,11 @@ GetPartitionTransformFromSpecField(Oid relationId, IcebergPartitionSpecField * s DataFileSchema *schema = GetDataFileSchemaForTable(relationId); - DataFileSchemaField *sourceField = GetDataFileSchemaFieldById(schema, specField->source_id); - - transform->sourceField = *sourceField; + transform->sourceField = GetDataFileSchemaFieldById(schema, specField->source_id); } /* parse transform name */ - ParseTransformName(transform->specField.transform, + ParseTransformName(transform->specField->transform, &transform->parsedTransform.type, &transform->parsedTransform.bucketCount, &transform->parsedTransform.truncateLen); @@ -415,8 +411,8 @@ ApplyPartitionTransformToTuple(IcebergPartitionTransform * transform, TupleTable { PartitionField *field = palloc0(sizeof(PartitionField)); - field->field_name = pstrdup(transform->specField.name); - field->field_id = transform->specField.field_id; + field->field_name = pstrdup(transform->specField->name); + field->field_id = transform->specField->field_id; bool isNull = false; Datum columnValue = slot_getattr(slot, transform->attnum, &isNull); @@ -455,7 +451,7 @@ ApplyPartitionTransformToTuple(IcebergPartitionTransform * transform, TupleTable ereport(ERROR, (errcode(ERRCODE_FEATURE_NOT_SUPPORTED), errmsg("applying transform %s is not yet support ", - transform->specField.transform))); + transform->specField->transform))); } field->value_type = GetTransformResultAvroType(transform); @@ -479,7 +475,7 @@ ApplyIdentityTransformToColumn(IcebergPartitionTransform * transform, Datum colu return NULL; } - return PGIcebergBinarySerializePartitionFieldValue(columnValue, transform->sourceField.type, + return PGIcebergBinarySerializePartitionFieldValue(columnValue, transform->sourceField->type, transform->pgType, valueSize); } @@ -840,7 +836,7 @@ ApplyBucketTransformToColumn(IcebergPartitionTransform * transform, Datum column else if (transform->pgType.postgresTypeOid == UUIDOID) { size_t valueSize = 0; - unsigned char *value = PGIcebergBinarySerializePartitionFieldValue(columnValue, transform->sourceField.type, + unsigned char *value = PGIcebergBinarySerializePartitionFieldValue(columnValue, transform->sourceField->type, transform->pgType, &valueSize); *bucketValue = (MurmurHash3_32_Bytes(value, valueSize) & INT32_MAX) % transform->parsedTransform.bucketCount; @@ -848,7 +844,7 @@ ApplyBucketTransformToColumn(IcebergPartitionTransform * transform, Datum column else if (transform->pgType.postgresTypeOid == NUMERICOID) { size_t valueSize = 0; - unsigned char *value = PGIcebergBinarySerializePartitionFieldValue(columnValue, transform->sourceField.type, + unsigned char *value = PGIcebergBinarySerializePartitionFieldValue(columnValue, transform->sourceField->type, transform->pgType, &valueSize); *bucketValue = (MurmurHash3_32_Bytes(value, valueSize) & INT32_MAX) % transform->parsedTransform.bucketCount; diff --git a/pg_lake_table/src/fdw/partitioning/partition_by_parser.c b/pg_lake_table/src/fdw/partitioning/partition_by_parser.c index 98be72835..0219da363 100644 --- a/pg_lake_table/src/fdw/partitioning/partition_by_parser.c +++ b/pg_lake_table/src/fdw/partitioning/partition_by_parser.c @@ -532,27 +532,28 @@ AnalyzeIcebergTablePartitionBy(Oid relationId, List *parsedTransforms) /* 2) Check scalar column */ EnsureTransformSourceColumnScalar(parsedTransform, sourceField); - analyzedTransform->sourceField = *sourceField; + analyzedTransform->sourceField = sourceField; /* set transform name */ - analyzedTransform->specField.transform = GenerateTransformName(analyzedTransform); - analyzedTransform->specField.transform_length = strlen(analyzedTransform->specField.transform); + analyzedTransform->specField = palloc0(sizeof(IcebergPartitionSpecField)); + analyzedTransform->specField->transform = GenerateTransformName(analyzedTransform); + analyzedTransform->specField->transform_length = strlen(analyzedTransform->specField->transform); /* set partition field name */ - analyzedTransform->specField.name = GeneratePartitionFieldName(analyzedTransform, relationId); - analyzedTransform->specField.name_length = strlen(analyzedTransform->specField.name); + analyzedTransform->specField->name = GeneratePartitionFieldName(analyzedTransform, relationId); + analyzedTransform->specField->name_length = strlen(analyzedTransform->specField->name); /* set partition field id */ - analyzedTransform->specField.field_id = ++largestPartitionFieldId; + analyzedTransform->specField->field_id = ++largestPartitionFieldId; /* set source field id */ - analyzedTransform->specField.source_id = sourceField->id; - analyzedTransform->specField.source_ids_length = 1; - analyzedTransform->specField.source_ids = palloc0(sizeof(int) * analyzedTransform->specField.source_ids_length); - analyzedTransform->specField.source_ids[0] = analyzedTransform->specField.source_id; + analyzedTransform->specField->source_id = sourceField->id; + analyzedTransform->specField->source_ids_length = 1; + analyzedTransform->specField->source_ids = palloc0(sizeof(int) * analyzedTransform->specField->source_ids_length); + analyzedTransform->specField->source_ids[0] = analyzedTransform->specField->source_id; /* 3) Check column type compatibility. */ EnsureValidTypeForTransform(analyzedTransform->parsedTransform.type, analyzedTransform->pgType.postgresTypeOid); - + analyzedTransforms = lappend(analyzedTransforms, analyzedTransform); } @@ -832,11 +833,11 @@ EnsureNoDuplicateTransforms(List *transforms) IcebergPartitionTransform *other = list_nth(transforms, j); if (transform->parsedTransform.type == other->parsedTransform.type && - transform->sourceField.id == other->sourceField.id) + transform->sourceField->id == other->sourceField->id) ereport(ERROR, (errcode(ERRCODE_DUPLICATE_OBJECT), errmsg("\"%s\" transform on column \"%s\" appears multiple times in partition spec", - transform->specField.transform, transform->parsedTransform.columnName))); + transform->specField->transform, transform->parsedTransform.columnName))); } } }