szehon-ho commented on code in PR #17957:
URL: https://github.com/apache/iceberg/pull/17957#discussion_r4137377788
##########
spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/source/SparkTable.java:
##########
@@ -158,6 +173,89 @@ public Set<TableCapability> capabilities() {
return capabilities;
}
+ @Override
+ public boolean supportsColumnChange(TableChange.ColumnChange change) {
+ if (isMapKeyChange(change)) {
+ return false;
+ }
+
+ if (change instanceof TableChange.AddColumn) {
+ TableChange.AddColumn add = (TableChange.AddColumn) change;
+ return add.isNullable() && add.defaultValue() == null &&
canConvert(add.dataType());
+ } else if (change instanceof TableChange.UpdateColumnType) {
+ return supportsTypeUpdate((TableChange.UpdateColumnType) change);
+ } else if (change instanceof TableChange.UpdateColumnNullability) {
+ return ((TableChange.UpdateColumnNullability) change).nullable();
+ } else if (change instanceof TableChange.DeleteColumn) {
+ return supportsDeleteColumn((TableChange.DeleteColumn) change);
+ } else {
+ return change instanceof TableChange.RenameColumn
+ || change instanceof TableChange.UpdateColumnComment
+ || change instanceof TableChange.UpdateColumnPosition;
+ }
+ }
+
+ private static Set<Integer> mapKeyFieldIds(Schema schema) {
+ Set<Integer> keyFieldIds = Sets.newHashSet();
+ for (Types.NestedField field :
TypeUtil.indexById(schema.asStruct()).values()) {
+ if (field.type().isMapType()) {
+ Types.MapType map = field.type().asMapType();
+
keyFieldIds.addAll(TypeUtil.indexById(Types.StructType.of(map.fields().get(0))).keySet());
+ }
+ }
+
+ return keyFieldIds;
+ }
+
+ private boolean isMapKeyChange(TableChange.ColumnChange change) {
+ String[] fieldNames = change.fieldNames();
+ int pathLength =
+ change instanceof TableChange.AddColumn ? fieldNames.length - 1 :
fieldNames.length;
+ if (pathLength == 0) {
+ return false;
+ }
+
+ Types.NestedField field =
+ schema.findField(String.join(".", Arrays.copyOf(fieldNames,
pathLength)));
+ return field != null && mapKeyFieldIds.contains(field.fieldId());
+ }
+
+ private boolean supportsTypeUpdate(TableChange.UpdateColumnType update) {
+ Types.NestedField field = schema.findField(String.join(".",
update.fieldNames()));
+ if (field == null) {
+ return false;
+ }
+
+ Type newType = tryConvert(update.newDataType());
+ return newType != null
+ && newType.isPrimitiveType()
+ && TypeUtil.isPromotionAllowed(field.type(),
newType.asPrimitiveType());
Review Comment:
Could we reject type updates when the converted Iceberg type does not map
back to the requested Spark type, and add an INSERT or MERGE test? For an INT
target with a SMALLINT/TINYINT source, ShortType/ByteType converts to Iceberg
int, so `isPromotionAllowed(int, int)` returns true and `updateColumn` is a
no-op. Spark reloads the table, sees the same pending change, and fails with
`UNSUPPORTED_AUTO_SCHEMA_EVOLUTION_CHANGES` instead of casting the source to
INT.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]