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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@
import java.util.stream.Collectors;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.hive.conf.HiveConf;
import org.apache.hadoop.hive.ql.Context;
import org.apache.hadoop.hive.ql.Context.Operation;
import org.apache.hadoop.hive.ql.metadata.RowLineageUtils;
import org.apache.hadoop.hive.ql.security.authorization.HiveCustomStorageHandlerUtils;
import org.apache.hadoop.hive.ql.session.SessionStateUtil;
Expand Down Expand Up @@ -158,7 +158,8 @@ public void initialize(Configuration conf, Properties serDeProperties,
private static Schema projectedSchema(Configuration conf, Properties serDeProperties,
Schema tableSchema, Map<String, String> jobConf) {
String tableName = serDeProperties.getProperty(Catalogs.NAME);
Context.Operation operation = HiveCustomStorageHandlerUtils.getWriteOperation(conf::get, tableName);
Operation operation = HiveCustomStorageHandlerUtils.getWriteOperation(conf::get, tableName);
boolean copyOnWrite = HiveCustomStorageHandlerUtils.isCopyOnWrite(conf::get, tableName);

if (operation == null) {
jobConf.put(InputFormatConfig.CASE_SENSITIVE, "false");
Expand All @@ -179,8 +180,7 @@ private static Schema projectedSchema(Configuration conf, Properties serDeProper
return projectedSchema;
}
}
boolean isCOW = IcebergTableUtil.isCopyOnWriteMode(operation, conf::get);
if (isCOW) {
if (copyOnWrite) {
return getSchemaWithRowLineage(
IcebergAcidUtil.createSerdeSchemaForDelete(tableSchema.columns(), false), conf);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,7 @@ public class WriterBuilder {
private TaskAttemptID attemptID;
private String queryId;
private Operation operation;
private final boolean copyOnWrite;

// A task may write multiple output files using multiple writers. Each of them must have a unique operationId.
private static AtomicInteger operationNum = new AtomicInteger(0);
Expand All @@ -87,6 +88,7 @@ private WriterBuilder(Table table, UnaryOperator<String> ops) {
this.tableName = ops.apply(Catalogs.NAME);
this.context = new Context(table.properties(), ops, tableName);
this.operation = HiveCustomStorageHandlerUtils.getWriteOperation(ops, tableName);
this.copyOnWrite = HiveCustomStorageHandlerUtils.isCopyOnWrite(ops, tableName);
this.rewritableDeletes = () -> rewritableDeletes(ops);
}

Expand Down Expand Up @@ -136,9 +138,7 @@ public HiveIcebergWriter build() {
.build();

HiveIcebergWriter writer;
boolean isCOW = IcebergTableUtil.isCopyOnWriteMode(operation, table.properties()::getOrDefault);

if (isCOW) {
if (copyOnWrite) {
writer = new HiveIcebergCopyOnWriteRecordWriter(table, writerFactory, dataFileFactory, shouldAddRowLineageColumns,
context);
} else {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
CREATE TABLE merge_source (id INT, data STRING);
INSERT INTO merge_source VALUES (2, 'banana_source');

CREATE TABLE ice_cow_merge_delete_only (id INT, data STRING)
STORED BY ICEBERG
TBLPROPERTIES ('format-version'='3', 'write.delete.mode'='copy-on-write');

INSERT INTO ice_cow_merge_delete_only VALUES (1, 'apple'), (2, 'banana'), (3, 'cherry');

SELECT id, data FROM ice_cow_merge_delete_only ORDER BY id;

MERGE INTO ice_cow_merge_delete_only t
USING merge_source s
ON t.id = s.id
WHEN MATCHED THEN DELETE;

SELECT id, data FROM ice_cow_merge_delete_only ORDER BY id;
36 changes: 36 additions & 0 deletions iceberg/iceberg-handler/src/test/queries/positive/row_lineage.q
Original file line number Diff line number Diff line change
Expand Up @@ -134,3 +134,39 @@ DELETE FROM ice_cow_delete_part WHERE id = 2 OR id = 3;
SELECT id, data, part, ROW__LINEAGE__ID, LAST__UPDATED__SEQUENCE__NUMBER
FROM ice_cow_delete_part
ORDER BY id;

-- cow merge delete only
CREATE TABLE merge_source (
id INT,
data STRING
);

INSERT INTO merge_source VALUES
(2, 'banana_source');

CREATE TABLE ice_cow_merge_delete_only (
id INT,
data STRING
)
STORED BY iceberg
TBLPROPERTIES ('format-version'='3', 'write.delete.mode'='copy-on-write');

-- Snapshot 1: Sequence 1
INSERT INTO ice_cow_merge_delete_only VALUES
(1, 'apple'),
(2, 'banana'),
(3, 'cherry');

SELECT id, data, ROW__LINEAGE__ID, LAST__UPDATED__SEQUENCE__NUMBER
FROM ice_cow_merge_delete_only
ORDER BY id;

MERGE INTO ice_cow_merge_delete_only t
USING merge_source s
ON t.id = s.id
WHEN MATCHED THEN DELETE;

-- Verification: id=1 and id=3 should perfectly retain their original lineage
SELECT id, data, ROW__LINEAGE__ID, LAST__UPDATED__SEQUENCE__NUMBER
FROM ice_cow_merge_delete_only
ORDER BY id;
Original file line number Diff line number Diff line change
@@ -0,0 +1,78 @@
PREHOOK: query: CREATE TABLE merge_source (id INT, data STRING)
PREHOOK: type: CREATETABLE
PREHOOK: Output: database:default
PREHOOK: Output: default@merge_source
POSTHOOK: query: CREATE TABLE merge_source (id INT, data STRING)
POSTHOOK: type: CREATETABLE
POSTHOOK: Output: database:default
POSTHOOK: Output: default@merge_source
PREHOOK: query: INSERT INTO merge_source VALUES (2, 'banana_source')
PREHOOK: type: QUERY
PREHOOK: Input: _dummy_database@_dummy_table
PREHOOK: Output: default@merge_source
POSTHOOK: query: INSERT INTO merge_source VALUES (2, 'banana_source')
POSTHOOK: type: QUERY
POSTHOOK: Input: _dummy_database@_dummy_table
POSTHOOK: Output: default@merge_source
POSTHOOK: Lineage: merge_source.data SCRIPT []
POSTHOOK: Lineage: merge_source.id SCRIPT []
PREHOOK: query: CREATE TABLE ice_cow_merge_delete_only (id INT, data STRING)
STORED BY ICEBERG
TBLPROPERTIES ('format-version'='3', 'write.delete.mode'='copy-on-write')
PREHOOK: type: CREATETABLE
PREHOOK: Output: database:default
PREHOOK: Output: default@ice_cow_merge_delete_only
POSTHOOK: query: CREATE TABLE ice_cow_merge_delete_only (id INT, data STRING)
STORED BY ICEBERG
TBLPROPERTIES ('format-version'='3', 'write.delete.mode'='copy-on-write')
POSTHOOK: type: CREATETABLE
POSTHOOK: Output: database:default
POSTHOOK: Output: default@ice_cow_merge_delete_only
PREHOOK: query: INSERT INTO ice_cow_merge_delete_only VALUES (1, 'apple'), (2, 'banana'), (3, 'cherry')
PREHOOK: type: QUERY
PREHOOK: Input: _dummy_database@_dummy_table
PREHOOK: Output: default@ice_cow_merge_delete_only
POSTHOOK: query: INSERT INTO ice_cow_merge_delete_only VALUES (1, 'apple'), (2, 'banana'), (3, 'cherry')
POSTHOOK: type: QUERY
POSTHOOK: Input: _dummy_database@_dummy_table
POSTHOOK: Output: default@ice_cow_merge_delete_only
PREHOOK: query: SELECT id, data FROM ice_cow_merge_delete_only ORDER BY id
PREHOOK: type: QUERY
PREHOOK: Input: default@ice_cow_merge_delete_only
PREHOOK: Output: hdfs://### HDFS PATH ###
POSTHOOK: query: SELECT id, data FROM ice_cow_merge_delete_only ORDER BY id
POSTHOOK: type: QUERY
POSTHOOK: Input: default@ice_cow_merge_delete_only
POSTHOOK: Output: hdfs://### HDFS PATH ###
1 apple
2 banana
3 cherry
PREHOOK: query: MERGE INTO ice_cow_merge_delete_only t
USING merge_source s
ON t.id = s.id
WHEN MATCHED THEN DELETE
PREHOOK: type: QUERY
PREHOOK: Input: default@ice_cow_merge_delete_only
PREHOOK: Input: default@merge_source
PREHOOK: Output: default@ice_cow_merge_delete_only
PREHOOK: Output: default@merge_tmp_table
POSTHOOK: query: MERGE INTO ice_cow_merge_delete_only t
USING merge_source s
ON t.id = s.id
WHEN MATCHED THEN DELETE
POSTHOOK: type: QUERY
POSTHOOK: Input: default@ice_cow_merge_delete_only
POSTHOOK: Input: default@merge_source
POSTHOOK: Output: default@ice_cow_merge_delete_only
POSTHOOK: Output: default@merge_tmp_table
POSTHOOK: Lineage: merge_tmp_table.val EXPRESSION [(ice_cow_merge_delete_only)ice_cow_merge_delete_only.null, ]
PREHOOK: query: SELECT id, data FROM ice_cow_merge_delete_only ORDER BY id
PREHOOK: type: QUERY
PREHOOK: Input: default@ice_cow_merge_delete_only
PREHOOK: Output: hdfs://### HDFS PATH ###
POSTHOOK: query: SELECT id, data FROM ice_cow_merge_delete_only ORDER BY id
POSTHOOK: type: QUERY
POSTHOOK: Input: default@ice_cow_merge_delete_only
POSTHOOK: Output: hdfs://### HDFS PATH ###
1 apple
3 cherry
106 changes: 106 additions & 0 deletions iceberg/iceberg-handler/src/test/results/positive/row_lineage.q.out
Original file line number Diff line number Diff line change
Expand Up @@ -469,3 +469,109 @@ POSTHOOK: Output: hdfs://### HDFS PATH ###
1 apple p1 0 1
4 date p2 3 1
5 elderberry p1 4 2
PREHOOK: query: CREATE TABLE merge_source (
id INT,
data STRING
)
PREHOOK: type: CREATETABLE
PREHOOK: Output: database:default
PREHOOK: Output: default@merge_source
POSTHOOK: query: CREATE TABLE merge_source (
id INT,
data STRING
)
POSTHOOK: type: CREATETABLE
POSTHOOK: Output: database:default
POSTHOOK: Output: default@merge_source
PREHOOK: query: INSERT INTO merge_source VALUES
(2, 'banana_source')
PREHOOK: type: QUERY
PREHOOK: Input: _dummy_database@_dummy_table
PREHOOK: Output: default@merge_source
POSTHOOK: query: INSERT INTO merge_source VALUES
(2, 'banana_source')
POSTHOOK: type: QUERY
POSTHOOK: Input: _dummy_database@_dummy_table
POSTHOOK: Output: default@merge_source
POSTHOOK: Lineage: merge_source.data SCRIPT []
POSTHOOK: Lineage: merge_source.id SCRIPT []
PREHOOK: query: CREATE TABLE ice_cow_merge_delete_only (
id INT,
data STRING
)
STORED BY iceberg
TBLPROPERTIES ('format-version'='3', 'write.delete.mode'='copy-on-write')
PREHOOK: type: CREATETABLE
PREHOOK: Output: database:default
PREHOOK: Output: default@ice_cow_merge_delete_only
POSTHOOK: query: CREATE TABLE ice_cow_merge_delete_only (
id INT,
data STRING
)
STORED BY iceberg
TBLPROPERTIES ('format-version'='3', 'write.delete.mode'='copy-on-write')
POSTHOOK: type: CREATETABLE
POSTHOOK: Output: database:default
POSTHOOK: Output: default@ice_cow_merge_delete_only
PREHOOK: query: INSERT INTO ice_cow_merge_delete_only VALUES
(1, 'apple'),
(2, 'banana'),
(3, 'cherry')
PREHOOK: type: QUERY
PREHOOK: Input: _dummy_database@_dummy_table
PREHOOK: Output: default@ice_cow_merge_delete_only
POSTHOOK: query: INSERT INTO ice_cow_merge_delete_only VALUES
(1, 'apple'),
(2, 'banana'),
(3, 'cherry')
POSTHOOK: type: QUERY
POSTHOOK: Input: _dummy_database@_dummy_table
POSTHOOK: Output: default@ice_cow_merge_delete_only
PREHOOK: query: SELECT id, data, ROW__LINEAGE__ID, LAST__UPDATED__SEQUENCE__NUMBER
FROM ice_cow_merge_delete_only
ORDER BY id
PREHOOK: type: QUERY
PREHOOK: Input: default@ice_cow_merge_delete_only
PREHOOK: Output: hdfs://### HDFS PATH ###
POSTHOOK: query: SELECT id, data, ROW__LINEAGE__ID, LAST__UPDATED__SEQUENCE__NUMBER
FROM ice_cow_merge_delete_only
ORDER BY id
POSTHOOK: type: QUERY
POSTHOOK: Input: default@ice_cow_merge_delete_only
POSTHOOK: Output: hdfs://### HDFS PATH ###
1 apple 0 1
2 banana 1 1
3 cherry 2 1
PREHOOK: query: MERGE INTO ice_cow_merge_delete_only t
USING merge_source s
ON t.id = s.id
WHEN MATCHED THEN DELETE
PREHOOK: type: QUERY
PREHOOK: Input: default@ice_cow_merge_delete_only
PREHOOK: Input: default@merge_source
PREHOOK: Output: default@ice_cow_merge_delete_only
PREHOOK: Output: default@merge_tmp_table
POSTHOOK: query: MERGE INTO ice_cow_merge_delete_only t
USING merge_source s
ON t.id = s.id
WHEN MATCHED THEN DELETE
POSTHOOK: type: QUERY
POSTHOOK: Input: default@ice_cow_merge_delete_only
POSTHOOK: Input: default@merge_source
POSTHOOK: Output: default@ice_cow_merge_delete_only
POSTHOOK: Output: default@merge_tmp_table
POSTHOOK: Lineage: merge_tmp_table.val EXPRESSION [(ice_cow_merge_delete_only)ice_cow_merge_delete_only.null, ]
PREHOOK: query: SELECT id, data, ROW__LINEAGE__ID, LAST__UPDATED__SEQUENCE__NUMBER
FROM ice_cow_merge_delete_only
ORDER BY id
PREHOOK: type: QUERY
PREHOOK: Input: default@ice_cow_merge_delete_only
PREHOOK: Output: hdfs://### HDFS PATH ###
POSTHOOK: query: SELECT id, data, ROW__LINEAGE__ID, LAST__UPDATED__SEQUENCE__NUMBER
FROM ice_cow_merge_delete_only
ORDER BY id
POSTHOOK: type: QUERY
POSTHOOK: Input: default@ice_cow_merge_delete_only
POSTHOOK: Output: hdfs://### HDFS PATH ###
1 apple 0 1
3 cherry 2 1
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@

import static org.apache.hadoop.hive.conf.HiveConf.ConfVars.HIVE_TEMPORARY_TABLE_STORAGE;
import static org.apache.hadoop.hive.ql.security.authorization.HiveCustomStorageHandlerUtils.MERGE_TASK_ENABLED;
import static org.apache.hadoop.hive.ql.security.authorization.HiveCustomStorageHandlerUtils.setCopyOnWrite;
import static org.apache.hadoop.hive.ql.security.authorization.HiveCustomStorageHandlerUtils.setMergeTaskEnabled;
import static org.apache.hadoop.hive.ql.security.authorization.HiveCustomStorageHandlerUtils.setWriteOperation;
import static org.apache.hadoop.hive.ql.security.authorization.HiveCustomStorageHandlerUtils.setWriteOperationIsSorted;
Expand Down Expand Up @@ -640,6 +641,7 @@ protected void initializeOp(Configuration hconf) throws HiveException {

jc = new JobConf(hconf);
setWriteOperation(jc, getConf().getTableInfo().getTableName(), getConf().getWriteOperation());
setCopyOnWrite(jc, getConf().getTableInfo().getTableName(), getConf().isCopyOnWrite());
setWriteOperationIsSorted(jc, getConf().getTableInfo().getTableName(),
dpCtx != null && dpCtx.hasCustomPartitionOrSortExpression());
setMergeTaskEnabled(jc, getConf().getTableInfo().getTableName(),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8661,6 +8661,13 @@ private FileSinkDesc createFileSinkDesc(String dest, TableDesc table_desc,
}

fileSinkDesc.setWriteOperation(writeOperation);
if (writeOperation != Context.Operation.OTHER
&& dest_tab != null
&& dest_tab.getStorageHandler() != null) {
boolean copyOnWrite =
dest_tab.getStorageHandler().shouldOverwrite(dest_tab, ctx.getOperation());
fileSinkDesc.setCopyOnWrite(copyOnWrite);
}

fileSinkDesc.setTemporary(destTableIsTemporary);
fileSinkDesc.setMaterialization(destTableIsMaterialization);
Expand Down
11 changes: 11 additions & 0 deletions ql/src/java/org/apache/hadoop/hive/ql/plan/FileSinkDesc.java
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,7 @@ public enum DPSortState {
private Path destPath;
private boolean isHiveServerQuery;
private boolean isMerge;
private boolean copyOnWrite = false;
private boolean isMmCtas;

private Set<FileStatus> filesToFetch = null;
Expand Down Expand Up @@ -199,6 +200,8 @@ public Object clone() throws CloneNotSupportedException {
ret.setStatsReliable(statsReliable);
ret.setDpSortState(dpSortState);
ret.setWriteType(writeType);
ret.setWriteOperation(writeOperation);
ret.setCopyOnWrite(copyOnWrite);
ret.setTableWriteId(tableWriteId);
ret.setStatementId(statementId);
ret.setStatsTmpDir(statsTmpDir);
Expand Down Expand Up @@ -686,6 +689,14 @@ public boolean isMmCtas() {
return isMmCtas;
}

public void setCopyOnWrite(boolean copyOnWrite) {
this.copyOnWrite = copyOnWrite;
}

public boolean isCopyOnWrite() {
return copyOnWrite;
}

@Explain(displayName = "bucketingVersion", explainLevels = { Level.EXTENDED })
public int getBucketingVersionForExplain() {
return getBucketingVersion();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ public class HiveCustomStorageHandlerUtils {

public static final String WRITE_OPERATION_CONFIG_PREFIX = "file.sink.write.operation.";
public static final String WRITE_OPERATION_IS_SORTED = "file.sink.write.operation.sorted.";
public static final String IS_COPY_ON_WRITE_CONFIG_PREFIX = "file.sink.is.copy.on.write.";

public static final String MERGE_TASK_ENABLED = "file.sink.merge.task.enabled.";

Expand Down Expand Up @@ -96,4 +97,16 @@ public static boolean isMergeTaskEnabled(UnaryOperator<String> ops, String table
String operation = ops.apply(MERGE_TASK_ENABLED + tableName);
return Boolean.parseBoolean(operation);
}

public static void setCopyOnWrite(Configuration conf, String tableName, boolean copyOnWrite) {
if (conf == null || tableName == null) {
return;
}
conf.setBoolean(IS_COPY_ON_WRITE_CONFIG_PREFIX + tableName, copyOnWrite);
}

public static boolean isCopyOnWrite(UnaryOperator<String> ops, String tableName) {
String isCopyOnWrite = ops.apply(IS_COPY_ON_WRITE_CONFIG_PREFIX + tableName);
return Boolean.parseBoolean(isCopyOnWrite);
}
}
Loading