dramaticlly commented on code in PR #13913:
URL: https://github.com/apache/iceberg/pull/13913#discussion_r2317235312
##########
spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/procedures/ExpireSnapshotsProcedure.java:
##########
@@ -104,26 +118,28 @@ public ProcedureParameter[] parameters() {
@Override
@SuppressWarnings("checkstyle:CyclomaticComplexity")
public Iterator<Scan> call(InternalRow args) {
- Identifier tableIdent = toIdentifier(args.getString(0),
PARAMETERS[0].name());
- Long olderThanMillis = args.isNullAt(1) ? null :
DateTimeUtil.microsToMillis(args.getLong(1));
- Integer retainLastNum = args.isNullAt(2) ? null : args.getInt(2);
- Integer maxConcurrentDeletes = args.isNullAt(3) ? null : args.getInt(3);
- Boolean streamResult = args.isNullAt(4) ? null : args.getBoolean(4);
- long[] snapshotIds = args.isNullAt(5) ? null :
args.getArray(5).toLongArray();
- Boolean cleanExpiredMetadata = args.isNullAt(6) ? null :
args.getBoolean(6);
+ ProcedureInput input = new ProcedureInput(spark(), tableCatalog(),
PARAMETERS, args);
+ Identifier tableIdent = toIdentifier(input.asString(TABLE_PARAM),
TABLE_PARAM.name());
Review Comment:
```suggestion
Identifier tableIdent = input.ident(TABLE_PARAM);
```
##########
spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/procedures/ExpireSnapshotsProcedure.java:
##########
@@ -104,26 +118,28 @@ public ProcedureParameter[] parameters() {
@Override
@SuppressWarnings("checkstyle:CyclomaticComplexity")
public Iterator<Scan> call(InternalRow args) {
- Identifier tableIdent = toIdentifier(args.getString(0),
PARAMETERS[0].name());
- Long olderThanMillis = args.isNullAt(1) ? null :
DateTimeUtil.microsToMillis(args.getLong(1));
- Integer retainLastNum = args.isNullAt(2) ? null : args.getInt(2);
- Integer maxConcurrentDeletes = args.isNullAt(3) ? null : args.getInt(3);
- Boolean streamResult = args.isNullAt(4) ? null : args.getBoolean(4);
- long[] snapshotIds = args.isNullAt(5) ? null :
args.getArray(5).toLongArray();
- Boolean cleanExpiredMetadata = args.isNullAt(6) ? null :
args.getBoolean(6);
+ ProcedureInput input = new ProcedureInput(spark(), tableCatalog(),
PARAMETERS, args);
+ Identifier tableIdent = toIdentifier(input.asString(TABLE_PARAM),
TABLE_PARAM.name());
+ Long olderThanMillis = input.asTimestampLong(OLDER_THAN_PARAM, null);
+ Integer retainLastNum = input.asInt(RETAIN_LAST_PARAM, null);
+ Integer maxConcurrentDeletes = input.asInt(MAX_CONCURRENT_DELETES_PARAM,
null);
+ Boolean streamResult = input.asBoolean(STREAM_RESULTS_PARAM, null);
Review Comment:
why don't we have
```java
boolean streamResult = input.asBoolean(STREAM_RESULTS_PARAM, false);
```
and only set if it's true below? Same apply to cleanExpiredMetadata
##########
spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/procedures/SetCurrentSnapshotProcedure.java:
##########
@@ -87,9 +90,10 @@ public ProcedureParameter[] parameters() {
@Override
public Iterator<Scan> call(InternalRow args) {
- Identifier tableIdent = toIdentifier(args.getString(0),
PARAMETERS[0].name());
- Long snapshotId = args.isNullAt(1) ? null : args.getLong(1);
- String ref = args.isNullAt(2) ? null : args.getString(2);
+ ProcedureInput input = new ProcedureInput(spark(), tableCatalog(),
PARAMETERS, args);
+ Identifier tableIdent = toIdentifier(input.asString(TABLE_PARAM),
TABLE_PARAM.name());
Review Comment:
```suggestion
Identifier tableIdent = input.ident(TABLE_PARAM);
```
##########
spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/procedures/CherrypickSnapshotProcedure.java:
##########
@@ -83,8 +85,10 @@ public ProcedureParameter[] parameters() {
@Override
public Iterator<Scan> call(InternalRow args) {
- Identifier tableIdent = toIdentifier(args.getString(0),
PARAMETERS[0].name());
- long snapshotId = args.getLong(1);
+ ProcedureInput input = new ProcedureInput(spark(), tableCatalog(),
PARAMETERS, args);
+
+ Identifier tableIdent = toIdentifier(input.asString(TABLE_PARAM),
TABLE_PARAM.name());
Review Comment:
```suggestion
Identifier tableIdent = input.ident(TABLE_PARAM);
```
##########
spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/procedures/ExpireSnapshotsProcedure.java:
##########
@@ -104,26 +118,28 @@ public ProcedureParameter[] parameters() {
@Override
@SuppressWarnings("checkstyle:CyclomaticComplexity")
public Iterator<Scan> call(InternalRow args) {
- Identifier tableIdent = toIdentifier(args.getString(0),
PARAMETERS[0].name());
- Long olderThanMillis = args.isNullAt(1) ? null :
DateTimeUtil.microsToMillis(args.getLong(1));
- Integer retainLastNum = args.isNullAt(2) ? null : args.getInt(2);
- Integer maxConcurrentDeletes = args.isNullAt(3) ? null : args.getInt(3);
- Boolean streamResult = args.isNullAt(4) ? null : args.getBoolean(4);
- long[] snapshotIds = args.isNullAt(5) ? null :
args.getArray(5).toLongArray();
- Boolean cleanExpiredMetadata = args.isNullAt(6) ? null :
args.getBoolean(6);
+ ProcedureInput input = new ProcedureInput(spark(), tableCatalog(),
PARAMETERS, args);
+ Identifier tableIdent = toIdentifier(input.asString(TABLE_PARAM),
TABLE_PARAM.name());
+ Long olderThanMillis = input.asTimestampLong(OLDER_THAN_PARAM, null);
+ Integer retainLastNum = input.asInt(RETAIN_LAST_PARAM, null);
+ Integer maxConcurrentDeletes = input.asInt(MAX_CONCURRENT_DELETES_PARAM,
null);
+ Boolean streamResult = input.asBoolean(STREAM_RESULTS_PARAM, null);
+ Long[] snapshotIds = input.asLongArray(SNAPSHOT_IDS_PARAM, null);
+ Boolean cleanExpiredMetadata =
input.asBoolean(CLEAN_EXPIRED_METADATA_PARAM, null);
Preconditions.checkArgument(
maxConcurrentDeletes == null || maxConcurrentDeletes > 0,
"max_concurrent_deletes should have value > 0, value: %s",
maxConcurrentDeletes);
+ Long finalOlderThanMillis = olderThanMillis;
Review Comment:
do we need this finalOlderThanMillis? It seems we can run
`TestExpireSnapshotsProcedure` successfully without this
##########
spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/procedures/FastForwardBranchProcedure.java:
##########
@@ -75,9 +78,11 @@ public ProcedureParameter[] parameters() {
@Override
public Iterator<Scan> call(InternalRow args) {
- Identifier tableIdent = toIdentifier(args.getString(0),
PARAMETERS[0].name());
- String from = args.getString(1);
- String to = args.getString(2);
+ ProcedureInput input = new ProcedureInput(spark(), tableCatalog(),
PARAMETERS, args);
+
+ Identifier tableIdent = toIdentifier(input.asString(TABLE_PARAM),
TABLE_PARAM.name());
Review Comment:
```suggestion
Identifier tableIdent = input.ident(TABLE_PARAM);
```
##########
spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/procedures/RemoveOrphanFilesProcedure.java:
##########
@@ -105,58 +124,39 @@ public ProcedureParameter[] parameters() {
@Override
@SuppressWarnings("checkstyle:CyclomaticComplexity")
public Iterator<Scan> call(InternalRow args) {
- Identifier tableIdent = toIdentifier(args.getString(0),
PARAMETERS[0].name());
- Long olderThanMillis = args.isNullAt(1) ? null :
DateTimeUtil.microsToMillis(args.getLong(1));
- String location = args.isNullAt(2) ? null : args.getString(2);
- boolean dryRun = args.isNullAt(3) ? false : args.getBoolean(3);
- Integer maxConcurrentDeletes = args.isNullAt(4) ? null : args.getInt(4);
- String fileListView = args.isNullAt(5) ? null : args.getString(5);
+ ProcedureInput input = new ProcedureInput(spark(), tableCatalog(),
PARAMETERS, args);
+ Identifier tableIdent = toIdentifier(input.asString(TABLE_PARAM),
TABLE_PARAM.name());
Review Comment:
```suggestion
Identifier tableIdent = input.ident(TABLE_PARAM);
```
##########
spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/procedures/RollbackToTimestampProcedure.java:
##########
@@ -83,9 +84,10 @@ public ProcedureParameter[] parameters() {
@Override
public Iterator<Scan> call(InternalRow args) {
- Identifier tableIdent = toIdentifier(args.getString(0),
PARAMETERS[0].name());
+ ProcedureInput input = new ProcedureInput(spark(), tableCatalog(),
PARAMETERS, args);
+ Identifier tableIdent = toIdentifier(input.asString(TABLE_PARAM),
TABLE_PARAM.name());
Review Comment:
```suggestion
Identifier tableIdent = input.ident(TABLE_PARAM);
```
##########
spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/procedures/PublishChangesProcedure.java:
##########
@@ -87,8 +89,10 @@ public ProcedureParameter[] parameters() {
@Override
public Iterator<Scan> call(InternalRow args) {
- Identifier tableIdent = toIdentifier(args.getString(0),
PARAMETERS[0].name());
- String wapId = args.getString(1);
+ ProcedureInput input = new ProcedureInput(spark(), tableCatalog(),
PARAMETERS, args);
+
+ Identifier tableIdent = toIdentifier(input.asString(TABLE_PARAM),
TABLE_PARAM.name());
Review Comment:
```suggestion
Identifier tableIdent = input.ident(TABLE_PARAM);
```
##########
spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/procedures/ProcedureInput.java:
##########
@@ -119,6 +148,15 @@ public String[] asStringArray(ProcedureParameter param,
String[] defaultValue) {
defaultValue);
}
+ public DeleteOrphanFiles.PrefixMismatchMode
asPrefixMismatchMode(ProcedureParameter param) {
+ String modeAsString = asString(param, null);
+ DeleteOrphanFiles.PrefixMismatchMode prefixMismatchMode =
+ (modeAsString == null)
+ ? null
+ : DeleteOrphanFiles.PrefixMismatchMode.fromString(modeAsString);
+ return prefixMismatchMode;
Review Comment:
```suggestion
return (modeAsString == null)
? null
: DeleteOrphanFiles.PrefixMismatchMode.fromString(modeAsString);
```
##########
spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/procedures/RewriteManifestsProcedure.java:
##########
@@ -89,9 +92,10 @@ public ProcedureParameter[] parameters() {
@Override
public Iterator<Scan> call(InternalRow args) {
- Identifier tableIdent = toIdentifier(args.getString(0),
PARAMETERS[0].name());
- Boolean useCaching = args.isNullAt(1) ? null : args.getBoolean(1);
- Integer specId = args.isNullAt(2) ? null : args.getInt(2);
+ ProcedureInput input = new ProcedureInput(spark(), tableCatalog(),
PARAMETERS, args);
+ Identifier tableIdent = toIdentifier(input.asString(TABLE_PARAM),
TABLE_PARAM.name());
+ Boolean useCaching = input.asBoolean(USE_CACHING_PARAM, false);
Review Comment:
let's use primitive boolean as well
##########
spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/procedures/RegisterTableProcedure.java:
##########
@@ -81,9 +83,11 @@ public ProcedureParameter[] parameters() {
@Override
public Iterator<Scan> call(InternalRow args) {
+ ProcedureInput input = new ProcedureInput(spark(), tableCatalog(),
PARAMETERS, args);
TableIdentifier tableName =
- Spark3Util.identifierToTableIdentifier(toIdentifier(args.getString(0),
"table"));
- String metadataFile = args.getString(1);
+ Spark3Util.identifierToTableIdentifier(
+ toIdentifier(input.asString(TABLE_PARAM), TABLE_PARAM.name()));
Review Comment:
```suggestion
input.ident(TABLE_PARAM);
```
##########
spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/procedures/RollbackToSnapshotProcedure.java:
##########
@@ -82,8 +84,9 @@ public ProcedureParameter[] parameters() {
@Override
public Iterator<Scan> call(InternalRow args) {
- Identifier tableIdent = toIdentifier(args.getString(0),
PARAMETERS[0].name());
- long snapshotId = args.getLong(1);
+ ProcedureInput input = new ProcedureInput(spark(), tableCatalog(),
PARAMETERS, args);
+ Identifier tableIdent = toIdentifier(input.asString(TABLE_PARAM),
TABLE_PARAM.name());
Review Comment:
```suggestion
Identifier tableIdent = input.ident(TABLE_PARAM);
```
##########
spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/procedures/ProcedureInput.java:
##########
@@ -92,6 +110,17 @@ public Long asLong(ProcedureParameter param, Long
defaultValue) {
return args.isNullAt(ordinal) ? defaultValue : (Long)
args.getLong(ordinal);
}
+ public Long[] asLongArray(ProcedureParameter param) {
+ Long[] value = asLongArray(param, null);
+ Preconditions.checkArgument(value != null, "Parameter '%s' is not set",
param.name());
+ return value;
+ }
Review Comment:
looks like this is not needed for this change, let's add later when the
needs come
--
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]