diff --git a/examples/cli/src/main/java/io/unitycatalog/cli/TableCli.java b/examples/cli/src/main/java/io/unitycatalog/cli/TableCli.java index dcf876796d..97ad1846b0 100644 --- a/examples/cli/src/main/java/io/unitycatalog/cli/TableCli.java +++ b/examples/cli/src/main/java/io/unitycatalog/cli/TableCli.java @@ -57,7 +57,7 @@ public static void handle(CommandLine cmd, ApiClient apiClient, String authToken output = getTable(tablesApi, json); break; case CliUtils.READ: - output = readTable(temporaryCredentialsApi, tablesApi, json); + output = readTable(temporaryCredentialsApi, tablesApi, json, cmd); break; case CliUtils.WRITE: output = writeTable(temporaryCredentialsApi, tablesApi, json); @@ -220,8 +220,11 @@ private static String getTable(TablesApi tablesApi, JSONObject json) } private static String readTable( - TemporaryCredentialsApi temporaryCredentialsApi, TablesApi tablesApi, JSONObject json) - throws ApiException { + TemporaryCredentialsApi temporaryCredentialsApi, + TablesApi tablesApi, + JSONObject json, + CommandLine cmd) + throws ApiException, JsonProcessingException { String fullTableName = json.getString(CliParams.FULL_NAME.getServerParam()); TableInfo info = tablesApi.getTable( @@ -242,6 +245,15 @@ private static String readTable( new GenerateTemporaryTableCredential() .tableId(tableId) .operation(TableOperation.READ)); + boolean jsonOutput = + cmd.hasOption(CliUtils.OUTPUT) + && ("json".equals(cmd.getOptionValue(CliUtils.OUTPUT)) + || "jsonPretty".equals(cmd.getOptionValue(CliUtils.OUTPUT))); + if (jsonOutput) { + return objectWriter.writeValueAsString( + DeltaKernelUtils.readDeltaTableAsRecords( + info.getStorageLocation(), temporaryCredentials, maxResults)); + } return DeltaKernelUtils.readDeltaTable( info.getStorageLocation(), temporaryCredentials, maxResults); } catch (Exception e) { diff --git a/examples/cli/src/main/java/io/unitycatalog/cli/delta/DeltaKernelUtils.java b/examples/cli/src/main/java/io/unitycatalog/cli/delta/DeltaKernelUtils.java index 353f7bfd23..6244e3ae0a 100644 --- a/examples/cli/src/main/java/io/unitycatalog/cli/delta/DeltaKernelUtils.java +++ b/examples/cli/src/main/java/io/unitycatalog/cli/delta/DeltaKernelUtils.java @@ -25,7 +25,9 @@ import io.unitycatalog.client.model.ColumnInfo; import io.unitycatalog.client.model.TemporaryCredentials; import java.net.URI; +import java.util.ArrayList; import java.util.HashMap; +import java.util.LinkedHashMap; import java.util.List; import java.util.Locale; import java.util.Map; @@ -209,6 +211,32 @@ public static String readDeltaTable( } } + public static List> readDeltaTableAsRecords( + String tablePath, TemporaryCredentials temporaryCredentials, int maxResults) { + Engine engine = getEngine(URI.create(tablePath), temporaryCredentials); + try { + Table table = Table.forPath(engine, substituteSchemeForS3(tablePath)); + Snapshot snapshot = table.getLatestSnapshot(engine); + StructType readSchema = snapshot.getSchema(); + ScanBuilder scanBuilder = snapshot.getScanBuilder().withReadSchema(readSchema); + List rowData = + DeltaKernelReadUtils.readData(engine, readSchema, scanBuilder.build(), maxResults); + + List> records = new ArrayList<>(); + for (Row row : rowData) { + Map record = new LinkedHashMap<>(); + for (int colOrdinal = 0; colOrdinal < readSchema.length(); colOrdinal++) { + String colName = readSchema.at(colOrdinal).getName(); + record.put(colName, DeltaKernelReadUtils.getValue(row, colOrdinal)); + } + records.add(record); + } + return records; + } catch (Exception e) { + throw new IllegalArgumentException("Failed to read Delta table", e); + } + } + // TODO : INTERVAL, CHAR and NULL, ARRAY, MAP, STRUCT public static StructType getSchema(List columns) { StructType structType = new StructType();