Skip to content

Commit

Permalink
Flink: Add engine name, version to EnvironmentContext (apache#6184)
Browse files Browse the repository at this point in the history
  • Loading branch information
nastra authored Nov 14, 2022
1 parent abf3156 commit a9d3630
Show file tree
Hide file tree
Showing 3 changed files with 12 additions and 0 deletions.
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@
import org.apache.flink.util.StringUtils;
import org.apache.iceberg.CachingCatalog;
import org.apache.iceberg.DataFile;
import org.apache.iceberg.EnvironmentContext;
import org.apache.iceberg.FileScanTask;
import org.apache.iceberg.PartitionField;
import org.apache.iceberg.PartitionSpec;
Expand Down Expand Up @@ -113,6 +114,9 @@ public FlinkCatalog(
asNamespaceCatalog =
originalCatalog instanceof SupportsNamespaces ? (SupportsNamespaces) originalCatalog : null;
closeable = originalCatalog instanceof Closeable ? (Closeable) originalCatalog : null;

EnvironmentContext.put(EnvironmentContext.ENGINE_NAME, "flink");
EnvironmentContext.put(EnvironmentContext.ENGINE_VERSION, "1.14");
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@
import org.apache.flink.util.StringUtils;
import org.apache.iceberg.CachingCatalog;
import org.apache.iceberg.DataFile;
import org.apache.iceberg.EnvironmentContext;
import org.apache.iceberg.FileScanTask;
import org.apache.iceberg.PartitionField;
import org.apache.iceberg.PartitionSpec;
Expand Down Expand Up @@ -113,6 +114,9 @@ public FlinkCatalog(
asNamespaceCatalog =
originalCatalog instanceof SupportsNamespaces ? (SupportsNamespaces) originalCatalog : null;
closeable = originalCatalog instanceof Closeable ? (Closeable) originalCatalog : null;

EnvironmentContext.put(EnvironmentContext.ENGINE_NAME, "flink");
EnvironmentContext.put(EnvironmentContext.ENGINE_VERSION, "1.15");
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@
import org.apache.flink.util.StringUtils;
import org.apache.iceberg.CachingCatalog;
import org.apache.iceberg.DataFile;
import org.apache.iceberg.EnvironmentContext;
import org.apache.iceberg.FileScanTask;
import org.apache.iceberg.PartitionField;
import org.apache.iceberg.PartitionSpec;
Expand Down Expand Up @@ -113,6 +114,9 @@ public FlinkCatalog(
asNamespaceCatalog =
originalCatalog instanceof SupportsNamespaces ? (SupportsNamespaces) originalCatalog : null;
closeable = originalCatalog instanceof Closeable ? (Closeable) originalCatalog : null;

EnvironmentContext.put(EnvironmentContext.ENGINE_NAME, "flink");
EnvironmentContext.put(EnvironmentContext.ENGINE_VERSION, "1.16");
}

@Override
Expand Down

0 comments on commit a9d3630

Please sign in to comment.