Skip to content

Commit c871adc

Browse files
committed
Narrow only the write relation for column updates and address review comments
Rework column-level UPDATE narrowing so that RewriteUpdateTable keeps master's read relation and plan builders and narrows only the write relation to requiredDataAttributes(). The rewrite orders the write query so the columns the write reads come first, and the ColumnPruning case for RowLevelWrite prunes only the columns after the last one the write reads. - Rename DataWriter.writeUpdate to writeColumnUpdate and LogicalWriteInfo.updateSchema to columnUpdateSchema, and the related error to DATA_SOURCE_WRITE_COLUMN_UPDATE_NOT_IMPLEMENTED. - Reject a write whose distribution or ordering references a column the column update does not read (COLUMN_UPDATE_UNDECLARED_WRITE_REQUIREMENT_COLUMNS), and a declared column that does not exist (COLUMN_UPDATE_UNKNOWN_REQUIRED_DATA_ATTRIBUTE). - Plan copy-on-write scans from the columns the write query reads for writes that deliver narrow rows; other writes read every column as before. - Build the runtime group filter attribute map from the original table. - Mark the new APIs @SInCE 4.4.0 and move the MiMa exclude to 4.4. - Fix Javadoc and error messages, and add tests for the rows and metadata the connector receives, cross-column assignments, case-sensitive analysis and MERGE updatedColumns().
1 parent 5b2369f commit c871adc

32 files changed

Lines changed: 2110 additions & 1179 deletions

‎common/utils/src/main/resources/error/error-conditions.json‎

Lines changed: 22 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1260,6 +1260,12 @@
12601260
],
12611261
"sqlState" : "42000"
12621262
},
1263+
"COLUMN_UPDATE_METADATA_REQUIRED_DATA_ATTRIBUTE" : {
1264+
"message" : [
1265+
"Connector <connector> mixes in `SupportsColumnUpdates` but declared metadata column(s) <metadataAttributes> in `requiredDataAttributes()`. Declare metadata columns in `requiredMetadataAttributes()` instead."
1266+
],
1267+
"sqlState" : "42000"
1268+
},
12631269
"COLUMN_UPDATE_NESTED_REQUIRED_DATA_ATTRIBUTE" : {
12641270
"message" : [
12651271
"Connector <connector> mixes in `SupportsColumnUpdates` but declared nested field(s) <nestedAttributes> in `requiredDataAttributes()`. Declare the root struct column instead of a nested field."
@@ -1274,16 +1280,28 @@
12741280
},
12751281
"COLUMN_UPDATE_SPLIT_ROW_ID_NOT_DECLARED" : {
12761282
"message" : [
1277-
"Data source <connector> represents UPDATE as delete + insert but its `requiredDataAttributes()` does not include row ID column(s) <rowIds>. The REINSERT payload is projected down to `requiredDataAttributes()`, so an undeclared row ID column leaves the reinserted row with no identity for the connector to place it by. Include the row ID column(s) in `requiredDataAttributes()`."
1283+
"Connector <connector> mixes in `SupportsColumnUpdates` and represents UPDATE as delete and insert, but row ID column(s) <rowIds> do not reach the reinserted row, which then has no identity for the connector to place it by. Declare each data row ID column in `requiredDataAttributes()`, and return each metadata row ID column from `requiredMetadataAttributes()` with `MetadataColumn.PRESERVE_ON_REINSERT` set."
12781284
],
12791285
"sqlState" : "42000"
12801286
},
12811287
"COLUMN_UPDATE_SPLIT_ROW_ID_REASSIGNMENT" : {
12821288
"message" : [
1283-
"Cannot reassign row ID column(s) <rowIds> in a column-level UPDATE when data source <connector> represents UPDATE as delete + insert. The REINSERT path has no row-ID channel to reconstruct columns outside `requiredDataAttributes()`, so undeclared columns would become NULL. Either avoid reassigning row ID columns, or include every table column in `requiredDataAttributes()`."
1289+
"Connector <connector> mixes in `SupportsColumnUpdates` and represents UPDATE as delete and insert, so UPDATE cannot assign row ID column(s) <rowIds>. The reinserted row carries only the columns in `requiredDataAttributes()`, and with a new row ID it cannot be matched to the row whose other columns the connector must preserve. Do not assign row ID columns, or include every table column in `requiredDataAttributes()`."
1290+
],
1291+
"sqlState" : "42000"
1292+
},
1293+
"COLUMN_UPDATE_UNDECLARED_WRITE_REQUIREMENT_COLUMNS" : {
1294+
"message" : [
1295+
"Connector <connector> mixes in `SupportsColumnUpdates`, but its write requires a distribution or ordering by column(s) <columns> that the column-level UPDATE does not read. Declare each data column in `requiredDataAttributes()` and each metadata column in `requiredMetadataAttributes()`."
12841296
],
12851297
"sqlState" : "42000"
12861298
},
1299+
"COLUMN_UPDATE_UNKNOWN_REQUIRED_DATA_ATTRIBUTE" : {
1300+
"message" : [
1301+
"Connector <connector> mixes in `SupportsColumnUpdates` but declared column(s) <unknownAttributes> in `requiredDataAttributes()` that do not exist in the table."
1302+
],
1303+
"sqlState" : "42703"
1304+
},
12871305
"COMPARATOR_RETURNS_NULL" : {
12881306
"message" : [
12891307
"The comparator has returned a NULL for a comparison between <firstValue> and <secondValue>.",
@@ -2433,9 +2451,9 @@
24332451
],
24342452
"sqlState" : "42K03"
24352453
},
2436-
"DATA_SOURCE_WRITE_UPDATE_NOT_IMPLEMENTED" : {
2454+
"DATA_SOURCE_WRITE_COLUMN_UPDATE_NOT_IMPLEMENTED" : {
24372455
"message" : [
2438-
"<class> mixes in `SupportsColumnUpdates` but does not implement `writeUpdate`. Connectors that opt into narrow column-level updates must override `writeUpdate` to handle rows in the `LogicalWriteInfo.updateSchema()` layout; the default implementation would forward narrow rows to `write`, which expects the full `LogicalWriteInfo.schema()` layout."
2456+
"<class> does not override `writeColumnUpdate(record)`. A data writer that receives rows in the `LogicalWriteInfo.columnUpdateSchema()` layout must override `writeColumnUpdate(record)`, or `writeColumnUpdate(metadata, record)` if the operation returns metadata columns from `requiredMetadataAttributes()`."
24392457
],
24402458
"sqlState" : "0A000"
24412459
},

‎project/MimaExcludes.scala‎

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -46,7 +46,10 @@ object MimaExcludes {
4646
"org.apache.spark.ml.regression.DecisionTreeRegressionModel.numLeave"),
4747
// [SPARK-59154] Remove unused prediction variance helper after inlining its implementation.
4848
ProblemFilters.exclude[DirectMissingMethodProblem](
49-
"org.apache.spark.ml.regression.DecisionTreeRegressionModel.predictVariance")
49+
"org.apache.spark.ml.regression.DecisionTreeRegressionModel.predictVariance"),
50+
// [SPARK-58111] Write schema narrowing for column-level UPDATE in DSv2
51+
ProblemFilters.exclude[ReversedMissingMethodProblem](
52+
"org.apache.spark.sql.connector.write.RowLevelOperationInfo.updatedColumns")
5053
)
5154

5255
// Exclude rules for 4.3.x from 4.2.0 (add 4.3-specific filters below as needed).
@@ -73,10 +76,7 @@ object MimaExcludes {
7376
// [SPARK-57987] Add desc field to the SQL REST API Node case class
7477
ProblemFilters.exclude[DirectMissingMethodProblem]("org.apache.spark.status.api.v1.sql.Node.apply"),
7578
ProblemFilters.exclude[DirectMissingMethodProblem]("org.apache.spark.status.api.v1.sql.Node.copy"),
76-
ProblemFilters.exclude[MissingTypesProblem]("org.apache.spark.status.api.v1.sql.Node$"),
77-
// [SPARK-58111] Add scan and write schema narrowing for column-level UPDATEs in DSv2
78-
ProblemFilters.exclude[ReversedMissingMethodProblem](
79-
"org.apache.spark.sql.connector.write.RowLevelOperationInfo.updatedColumns")
79+
ProblemFilters.exclude[MissingTypesProblem]("org.apache.spark.status.api.v1.sql.Node$")
8080
)
8181

8282
// Exclude rules for 4.2.x from 4.1.0

‎sql/catalyst/src/main/java/org/apache/spark/sql/connector/write/DataWriter.java‎

Lines changed: 30 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -85,50 +85,56 @@ default void write(T metadata, T record) throws IOException {
8585
}
8686

8787
/**
88-
* Writes one updated, copied, or reinserted record with metadata.
88+
* Writes one updated or copied record with metadata in the column update layout.
8989
* <p>
90-
* Connectors that mix in {@link SupportsColumnUpdates} without also mixing in
91-
* {@link SupportsDelta} (i.e. group-based operations) receive records here in the schema
92-
* declared by {@link LogicalWriteInfo#updateSchema()}. An operation that mixes in both
93-
* instead delivers narrow rows through {@link DeltaWriter#update} / {@link DeltaWriter#reinsert}
94-
* and never calls this method. By default, delegates to {@link #writeUpdate(Object)} for
95-
* connectors that do not need metadata; implementations that do need it should override this
96-
* method directly.
90+
* When {@link LogicalWriteInfo#columnUpdateSchema()} is present for a row-level operation that
91+
* does not mix in {@link SupportsDelta}, Spark passes updated and copied records to this method
92+
* instead of {@link #write(Object, Object)} if the operation returns a non-empty
93+
* {@link RowLevelOperation#requiredMetadataAttributes()}, and to
94+
* {@link #writeColumnUpdate(Object)} otherwise. The record follows
95+
* {@link LogicalWriteInfo#columnUpdateSchema()} and the metadata follows
96+
* {@link LogicalWriteInfo#metadataSchema()}. Operations that mix in {@link SupportsDelta} receive
97+
* such rows through {@link DeltaWriter#update} and {@link DeltaWriter#reinsert} instead.
98+
* <p>
99+
* By default, delegates to {@link #writeColumnUpdate(Object)} and drops the metadata.
97100
* <p>
98101
* If this method fails (by throwing an exception), {@link #abort()} will be called and this
99102
* data writer is considered to have been failed.
100103
*
101104
* @throws IOException if failure happens during disk/network IO like writing files.
102-
* @throws SparkUnsupportedOperationException if the connector mixes in
103-
* {@link SupportsColumnUpdates} but overrides neither {@code writeUpdate} overload.
105+
* @throws SparkUnsupportedOperationException if neither this method nor
106+
* {@link #writeColumnUpdate(Object)} is overridden.
104107
*
105-
* @since 4.3.0
108+
* @since 4.4.0
106109
*/
107-
default void writeUpdate(T metadata, T record) throws IOException {
108-
writeUpdate(record);
110+
default void writeColumnUpdate(T metadata, T record) throws IOException {
111+
writeColumnUpdate(record);
109112
}
110113

111114
/**
112-
* Writes one updated, copied, or reinserted record without metadata.
115+
* Writes one updated or copied record without metadata in the column update layout.
116+
* <p>
117+
* When {@link LogicalWriteInfo#columnUpdateSchema()} is present for a row-level operation that
118+
* does not mix in {@link SupportsDelta}, Spark passes updated and copied records to this method
119+
* instead of {@link #write(Object)} if the operation returns no
120+
* {@link RowLevelOperation#requiredMetadataAttributes()}. The record follows
121+
* {@link LogicalWriteInfo#columnUpdateSchema()}.
113122
* <p>
114-
* Connectors that mix in {@link SupportsColumnUpdates} without also mixing in
115-
* {@link SupportsDelta} (i.e. group-based operations) receive records here in the schema
116-
* declared by {@link LogicalWriteInfo#updateSchema()}. Implementations must override this method,
117-
* or {@link #writeUpdate(Object, Object)}, when mixing in {@link SupportsColumnUpdates} without
118-
* {@link SupportsDelta}.
123+
* A writer for such an operation must override this method, unless the operation returns a
124+
* non-empty {@link RowLevelOperation#requiredMetadataAttributes()} and the writer overrides
125+
* {@link #writeColumnUpdate(Object, Object)}.
119126
* <p>
120127
* If this method fails (by throwing an exception), {@link #abort()} will be called and this
121128
* data writer is considered to have been failed.
122129
*
123130
* @throws IOException if failure happens during disk/network IO like writing files.
124-
* @throws SparkUnsupportedOperationException if the connector mixes in
125-
* {@link SupportsColumnUpdates} but overrides neither {@code writeUpdate} overload.
131+
* @throws SparkUnsupportedOperationException if this method is not overridden.
126132
*
127-
* @since 4.3.0
133+
* @since 4.4.0
128134
*/
129-
default void writeUpdate(T record) throws IOException {
135+
default void writeColumnUpdate(T record) throws IOException {
130136
throw new SparkUnsupportedOperationException(
131-
"DATA_SOURCE_WRITE_UPDATE_NOT_IMPLEMENTED",
137+
"DATA_SOURCE_WRITE_COLUMN_UPDATE_NOT_IMPLEMENTED",
132138
Map.of("class", getClass().getName()));
133139
}
134140

‎sql/catalyst/src/main/java/org/apache/spark/sql/connector/write/DeltaWriter.java‎

Lines changed: 4 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -41,10 +41,8 @@ public interface DeltaWriter<T> extends DataWriter<T> {
4141
/**
4242
* Updates a row.
4343
* <p>
44-
* When {@link LogicalWriteInfo#updateSchema()} is present, the {@code row} follows the narrow
45-
* layout it declares rather than the full table schema from {@link LogicalWriteInfo#schema()}.
46-
* It is present only for UPDATE on an operation that mixes in {@link SupportsColumnUpdates};
47-
* otherwise {@code row} follows the full table schema.
44+
* When {@link LogicalWriteInfo#columnUpdateSchema()} is present, the {@code row} follows it;
45+
* otherwise it follows {@link LogicalWriteInfo#schema()}.
4846
*
4947
* @param metadata values for metadata columns that were projected but are not part of the row ID
5048
* @param id a row ID to update
@@ -58,10 +56,8 @@ public interface DeltaWriter<T> extends DataWriter<T> {
5856
* <p>
5957
* This method handles the insert portion of updated rows split into deletes and inserts.
6058
* <p>
61-
* When {@link LogicalWriteInfo#updateSchema()} is present, the {@code row} follows the narrow
62-
* layout it declares rather than the full table schema from {@link LogicalWriteInfo#schema()}.
63-
* It is present only for UPDATE on an operation that mixes in {@link SupportsColumnUpdates};
64-
* otherwise {@code row} follows the full table schema.
59+
* When {@link LogicalWriteInfo#columnUpdateSchema()} is present, the {@code row} follows it;
60+
* otherwise it follows {@link LogicalWriteInfo#schema()}.
6561
*
6662
* @param metadata values for metadata columns
6763
* @param row a row to reinsert

‎sql/catalyst/src/main/java/org/apache/spark/sql/connector/write/LogicalWriteInfo.java‎

Lines changed: 12 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -47,6 +47,10 @@ public interface LogicalWriteInfo {
4747

4848
/**
4949
* the schema of the input data from Spark to data source.
50+
* <p>
51+
* When {@link #columnUpdateSchema()} is present, this schema covers only newly inserted rows, as
52+
* updated, copied, and reinserted rows follow {@link #columnUpdateSchema()}. It is then empty
53+
* when the command inserts no new rows, such as UPDATE.
5054
*/
5155
StructType schema();
5256

@@ -67,12 +71,16 @@ default Optional<StructType> metadataSchema() {
6771
}
6872

6973
/**
70-
* the narrow schema for updates. Present only when the connector mixes in
71-
* {@link SupportsColumnUpdates}.
74+
* the schema of updated, copied, and reinserted rows from Spark to data source in a
75+
* column-level update. Present when the operation mixes in {@link SupportsColumnUpdates} and
76+
* Spark delivers these rows with the columns of
77+
* {@link SupportsColumnUpdates#requiredDataAttributes()}, which currently happens only for
78+
* UPDATE. It covers every table column if every column is declared. When present,
79+
* {@link #schema()} covers only newly inserted rows.
7280
*
73-
* @since 4.3.0
81+
* @since 4.4.0
7482
*/
75-
default Optional<StructType> updateSchema() {
83+
default Optional<StructType> columnUpdateSchema() {
7684
return Optional.empty();
7785
}
7886
}

‎sql/catalyst/src/main/java/org/apache/spark/sql/connector/write/RowLevelOperationInfo.java‎

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -40,13 +40,14 @@ public interface RowLevelOperationInfo {
4040
Command command();
4141

4242
/**
43-
* Returns the columns being updated by this operation. Currently populated only for UPDATE;
44-
* DELETE and MERGE report an empty array.
43+
* Returns the columns being updated by this operation. Currently only UPDATE populates it;
44+
* other commands report an empty array.
4545
* <p>
46-
* Nested struct field updates are reported at root-column granularity
47-
* (e.g. {@code SET s.c1 = -1} returns {@code s}).
46+
* A column is reported only if it is assigned a new value, so identity assignments such as
47+
* {@code SET a = a} are excluded. Nested struct field updates are reported at root-column
48+
* granularity (e.g. {@code SET s.c1 = -1} returns {@code s}).
4849
*
49-
* @since 4.3.0
50+
* @since 4.4.0
5051
*/
5152
NamedReference[] updatedColumns();
5253
}

0 commit comments

Comments
 (0)