FrankChen021 commented on code in PR #19698:
URL: https://github.com/apache/druid/pull/19698#discussion_r3740843249
##########
server/src/main/java/org/apache/druid/metadata/SQLMetadataConnector.java:
##########
@@ -1054,6 +1061,196 @@ public void createAuditTable()
}
}
+ @Override
+ public void exportTable(
+ final String tableName,
+ final String outputPath
+ )
+ {
+ exportTable(tableName, outputPath, null);
+ }
+
+ /**
+ * Exports a table to a CSV file, emitting the given columns in the given
order.
+ *
+ * @param columns columns to export in the desired order, or null to export
all columns in the
+ * order reported by the database
+ */
+ public void exportTable(
+ final String tableName,
+ final String outputPath,
+ @Nullable final List<String> columns
+ )
+ {
+ exportTableWithJdbc(tableName, outputPath, columns);
+ }
+
+ /**
+ * Returns the columns of the given table, in the order reported by the
database.
+ * Returns an empty list if the table does not exist or the metadata cannot
be read.
+ *
+ * The lookup is scoped to the schema returned by {@link
#getMetadataTableSchema(Connection)}, which is the schema
+ * an unqualified table name resolves to. The table name is folded to the
case in which the database stores unquoted
+ * identifiers, since {@link DatabaseMetaData#getColumns} patterns are
case-sensitive.
+ */
+ public List<String> getTableColumns(final String tableName)
+ {
+ return getDBI().withHandle(handle -> {
+ final List<String> columns = new ArrayList<>();
+ try {
+ if (tableExists(handle, tableName)) {
+ final Connection conn = handle.getConnection();
+ final DatabaseMetaData dbMetaData = conn.getMetaData();
+ final String storedName = foldIdentifierCase(dbMetaData, tableName);
+ try (ResultSet rs = dbMetaData.getColumns(null,
getMetadataTableSchema(conn), storedName, null)) {
+ while (rs.next()) {
+ // '_' is a wildcard in the table name pattern, so match the
returned table name exactly
+ if (storedName.equals(rs.getString("TABLE_NAME"))) {
+ columns.add(rs.getString("COLUMN_NAME"));
+ }
+ }
+ }
+ }
+ }
+ catch (SQLException e) {
+ log.warn(e, "Could not read columns of table[%s].", tableName);
+ }
+ return columns;
+ });
+ }
+
+ /**
+ * Returns the schema that the Druid metadata tables live in, i.e. the
schema that an unqualified
+ * table name in a Druid SQL statement resolves to, or null if the schema is
unknown and lookups
+ * should not be scoped to a schema.
+ *
+ * Connectors that scope {@link #tableExists} to a configured schema must
override this so that
+ * both lookups agree.
+ */
+ @Nullable
+ protected String getMetadataTableSchema(final Connection connection) throws
SQLException
+ {
+ return connection.getSchema();
+ }
+
+ /**
+ * Folds the given identifier to the case in which the database stores
unquoted identifiers.
+ * {@link DatabaseMetaData} lookup patterns are case-sensitive, while an
unquoted identifier
+ * in a SQL statement is folded by the database (to lowercase in PostgreSQL,
to uppercase in
+ * Derby), so the folded form must be used to match the table a SQL
reference resolves to.
+ */
+ private static String foldIdentifierCase(
+ final DatabaseMetaData dbMetaData,
+ final String identifier
+ ) throws SQLException
+ {
+ if (dbMetaData.storesLowerCaseIdentifiers()) {
+ return StringUtils.toLowerCase(identifier);
+ }
+ if (dbMetaData.storesUpperCaseIdentifiers()) {
+ return StringUtils.toUpperCase(identifier);
+ }
+ return identifier;
+ }
+
+ /**
+ * Builds the select list for an export query, quoting each column with the
database's
+ * identifier quote string so that reserved words such as "end" are handled
correctly.
+ */
+ protected String makeExportSelectList(final Connection conn, final
List<String> columns) throws SQLException
+ {
+ final String quote = conn.getMetaData().getIdentifierQuoteString();
+ if (quote == null || " ".equals(quote)) {
+ return String.join(",", columns);
+ }
+ return columns.stream()
+ .map(column -> quote + StringUtils.replace(column, quote,
quote + quote) + quote)
+ .collect(Collectors.joining(","));
+ }
+
+ /**
+ * Exports a table to a CSV file using generic JDBC.
+ * Binary columns are hex-encoded and booleans are written as true/false
strings.
+ * Subclasses may override {@link #exportTable} with a database-specific
implementation
+ * while this method remains available for testing or fallback.
+ *
+ * @param columns columns to export in the desired order, or null to export
all columns
+ */
+ protected void exportTableWithJdbc(
+ final String tableName,
+ final String outputPath,
+ @Nullable final List<String> columns
+ )
+ {
+ // Use a transaction so that the connection has autoCommit=false.
+ // PostgreSQL JDBC requires autoCommit=false and a positive fetch size
+ // to use cursor-based streaming instead of buffering the entire ResultSet.
+ retryTransaction(
+ (TransactionCallback<Void>) (handle, status) -> {
+ final Connection conn = handle.getConnection();
+ final String selectList = columns == null || columns.isEmpty() ? "*"
: makeExportSelectList(conn, columns);
+ try (Statement stmt = conn.createStatement()) {
+ final int fetchSize = getStreamingFetchSize();
+ if (fetchSize > 0) {
+ stmt.setFetchSize(fetchSize);
+ }
+ try (ResultSet rs = stmt.executeQuery(StringUtils.format("SELECT
%s FROM %s", selectList, tableName));
Review Comment:
[P1] Export ignores configured PostgreSQL schema
With a non-public dbTableSchema, column discovery uses the configured schema
but this query remains unqualified. PostgreSQL may fail to find the table or
export a same-named table from another schema. Resolve and quote the configured
schema/table consistently.
##########
server/src/main/java/org/apache/druid/metadata/SQLMetadataConnector.java:
##########
@@ -1054,6 +1061,196 @@ public void createAuditTable()
}
}
+ @Override
+ public void exportTable(
+ final String tableName,
+ final String outputPath
+ )
+ {
+ exportTable(tableName, outputPath, null);
+ }
+
+ /**
+ * Exports a table to a CSV file, emitting the given columns in the given
order.
+ *
+ * @param columns columns to export in the desired order, or null to export
all columns in the
+ * order reported by the database
+ */
+ public void exportTable(
+ final String tableName,
+ final String outputPath,
+ @Nullable final List<String> columns
+ )
+ {
+ exportTableWithJdbc(tableName, outputPath, columns);
+ }
+
+ /**
+ * Returns the columns of the given table, in the order reported by the
database.
+ * Returns an empty list if the table does not exist or the metadata cannot
be read.
+ *
+ * The lookup is scoped to the schema returned by {@link
#getMetadataTableSchema(Connection)}, which is the schema
+ * an unqualified table name resolves to. The table name is folded to the
case in which the database stores unquoted
+ * identifiers, since {@link DatabaseMetaData#getColumns} patterns are
case-sensitive.
+ */
+ public List<String> getTableColumns(final String tableName)
+ {
+ return getDBI().withHandle(handle -> {
+ final List<String> columns = new ArrayList<>();
+ try {
+ if (tableExists(handle, tableName)) {
+ final Connection conn = handle.getConnection();
+ final DatabaseMetaData dbMetaData = conn.getMetaData();
+ final String storedName = foldIdentifierCase(dbMetaData, tableName);
+ try (ResultSet rs = dbMetaData.getColumns(null,
getMetadataTableSchema(conn), storedName, null)) {
+ while (rs.next()) {
+ // '_' is a wildcard in the table name pattern, so match the
returned table name exactly
+ if (storedName.equals(rs.getString("TABLE_NAME"))) {
+ columns.add(rs.getString("COLUMN_NAME"));
+ }
+ }
+ }
+ }
+ }
+ catch (SQLException e) {
+ log.warn(e, "Could not read columns of table[%s].", tableName);
+ }
+ return columns;
+ });
+ }
+
+ /**
+ * Returns the schema that the Druid metadata tables live in, i.e. the
schema that an unqualified
+ * table name in a Druid SQL statement resolves to, or null if the schema is
unknown and lookups
+ * should not be scoped to a schema.
+ *
+ * Connectors that scope {@link #tableExists} to a configured schema must
override this so that
+ * both lookups agree.
+ */
+ @Nullable
+ protected String getMetadataTableSchema(final Connection connection) throws
SQLException
+ {
+ return connection.getSchema();
+ }
+
+ /**
+ * Folds the given identifier to the case in which the database stores
unquoted identifiers.
+ * {@link DatabaseMetaData} lookup patterns are case-sensitive, while an
unquoted identifier
+ * in a SQL statement is folded by the database (to lowercase in PostgreSQL,
to uppercase in
+ * Derby), so the folded form must be used to match the table a SQL
reference resolves to.
+ */
+ private static String foldIdentifierCase(
+ final DatabaseMetaData dbMetaData,
+ final String identifier
+ ) throws SQLException
+ {
+ if (dbMetaData.storesLowerCaseIdentifiers()) {
+ return StringUtils.toLowerCase(identifier);
+ }
+ if (dbMetaData.storesUpperCaseIdentifiers()) {
+ return StringUtils.toUpperCase(identifier);
+ }
+ return identifier;
+ }
+
+ /**
+ * Builds the select list for an export query, quoting each column with the
database's
+ * identifier quote string so that reserved words such as "end" are handled
correctly.
+ */
+ protected String makeExportSelectList(final Connection conn, final
List<String> columns) throws SQLException
+ {
+ final String quote = conn.getMetaData().getIdentifierQuoteString();
+ if (quote == null || " ".equals(quote)) {
+ return String.join(",", columns);
+ }
+ return columns.stream()
+ .map(column -> quote + StringUtils.replace(column, quote,
quote + quote) + quote)
+ .collect(Collectors.joining(","));
+ }
+
+ /**
+ * Exports a table to a CSV file using generic JDBC.
+ * Binary columns are hex-encoded and booleans are written as true/false
strings.
+ * Subclasses may override {@link #exportTable} with a database-specific
implementation
+ * while this method remains available for testing or fallback.
+ *
+ * @param columns columns to export in the desired order, or null to export
all columns
+ */
+ protected void exportTableWithJdbc(
+ final String tableName,
+ final String outputPath,
+ @Nullable final List<String> columns
+ )
+ {
+ // Use a transaction so that the connection has autoCommit=false.
+ // PostgreSQL JDBC requires autoCommit=false and a positive fetch size
+ // to use cursor-based streaming instead of buffering the entire ResultSet.
+ retryTransaction(
+ (TransactionCallback<Void>) (handle, status) -> {
+ final Connection conn = handle.getConnection();
+ final String selectList = columns == null || columns.isEmpty() ? "*"
: makeExportSelectList(conn, columns);
+ try (Statement stmt = conn.createStatement()) {
+ final int fetchSize = getStreamingFetchSize();
+ if (fetchSize > 0) {
+ stmt.setFetchSize(fetchSize);
+ }
+ try (ResultSet rs = stmt.executeQuery(StringUtils.format("SELECT
%s FROM %s", selectList, tableName));
+ OutputStreamWriter writer =
+ new OutputStreamWriter(new FileOutputStream(outputPath),
StandardCharsets.UTF_8)) {
+ final int columnCount = rs.getMetaData().getColumnCount();
+ final List<String> values = new ArrayList<>(columnCount);
+ while (rs.next()) {
+ values.clear();
+ for (int i = 1; i <= columnCount; i++) {
+ values.add(readCsvValue(rs, i));
+ }
+ writer.write(String.join(",", values));
+ writer.write('\n');
+ }
+ }
+ }
+ return null;
+ },
+ QUIET_RETRIES,
+ DEFAULT_MAX_TRIES
+ );
+ }
+
+ /**
+ * Reads the given column of the current row as a CSV field. Binary values
are hex-encoded, booleans are written
+ * as true/false and NULLs as empty fields.
+ */
+ private static String readCsvValue(final ResultSet rs, final int column)
throws SQLException
+ {
+ final ResultSetMetaData meta = rs.getMetaData();
+ final int type = meta.getColumnType(column);
+ if (type == Types.BINARY || type == Types.VARBINARY || type ==
Types.LONGVARBINARY || type == Types.BLOB
+ || (type == Types.OTHER &&
"bytea".equalsIgnoreCase(meta.getColumnTypeName(column)))) {
+ final byte[] bytes = rs.getBytes(column);
+ return bytes == null ? "" : BaseEncoding.base16().encode(bytes);
+ } else if (type == Types.BOOLEAN || type == Types.BIT) {
+ final boolean value = rs.getBoolean(column);
+ return rs.wasNull() ? "" : String.valueOf(value);
+ } else {
+ return csvEscape(rs.getString(column));
Review Comment:
[P1] Nullable num_rows cannot be imported safely
A nullable num_rows is serialized as an empty field. The Derby and MySQL
commands documented for adding this column do not convert empty input to SQL
NULL: Derby rejects it and MySQL may coerce it to zero. Add equivalent NULL
handling for those import paths.
##########
server/src/main/java/org/apache/druid/metadata/SQLMetadataConnector.java:
##########
@@ -1054,6 +1061,196 @@ public void createAuditTable()
}
}
+ @Override
+ public void exportTable(
+ final String tableName,
+ final String outputPath
+ )
+ {
+ exportTable(tableName, outputPath, null);
+ }
+
+ /**
+ * Exports a table to a CSV file, emitting the given columns in the given
order.
+ *
+ * @param columns columns to export in the desired order, or null to export
all columns in the
+ * order reported by the database
+ */
+ public void exportTable(
+ final String tableName,
+ final String outputPath,
+ @Nullable final List<String> columns
+ )
+ {
+ exportTableWithJdbc(tableName, outputPath, columns);
+ }
+
+ /**
+ * Returns the columns of the given table, in the order reported by the
database.
+ * Returns an empty list if the table does not exist or the metadata cannot
be read.
+ *
+ * The lookup is scoped to the schema returned by {@link
#getMetadataTableSchema(Connection)}, which is the schema
+ * an unqualified table name resolves to. The table name is folded to the
case in which the database stores unquoted
+ * identifiers, since {@link DatabaseMetaData#getColumns} patterns are
case-sensitive.
+ */
+ public List<String> getTableColumns(final String tableName)
+ {
+ return getDBI().withHandle(handle -> {
+ final List<String> columns = new ArrayList<>();
+ try {
+ if (tableExists(handle, tableName)) {
+ final Connection conn = handle.getConnection();
+ final DatabaseMetaData dbMetaData = conn.getMetaData();
+ final String storedName = foldIdentifierCase(dbMetaData, tableName);
+ try (ResultSet rs = dbMetaData.getColumns(null,
getMetadataTableSchema(conn), storedName, null)) {
+ while (rs.next()) {
+ // '_' is a wildcard in the table name pattern, so match the
returned table name exactly
+ if (storedName.equals(rs.getString("TABLE_NAME"))) {
+ columns.add(rs.getString("COLUMN_NAME"));
+ }
+ }
+ }
+ }
+ }
+ catch (SQLException e) {
+ log.warn(e, "Could not read columns of table[%s].", tableName);
+ }
+ return columns;
+ });
+ }
+
+ /**
+ * Returns the schema that the Druid metadata tables live in, i.e. the
schema that an unqualified
+ * table name in a Druid SQL statement resolves to, or null if the schema is
unknown and lookups
+ * should not be scoped to a schema.
+ *
+ * Connectors that scope {@link #tableExists} to a configured schema must
override this so that
+ * both lookups agree.
+ */
+ @Nullable
+ protected String getMetadataTableSchema(final Connection connection) throws
SQLException
+ {
+ return connection.getSchema();
+ }
+
+ /**
+ * Folds the given identifier to the case in which the database stores
unquoted identifiers.
+ * {@link DatabaseMetaData} lookup patterns are case-sensitive, while an
unquoted identifier
+ * in a SQL statement is folded by the database (to lowercase in PostgreSQL,
to uppercase in
+ * Derby), so the folded form must be used to match the table a SQL
reference resolves to.
+ */
+ private static String foldIdentifierCase(
+ final DatabaseMetaData dbMetaData,
+ final String identifier
+ ) throws SQLException
+ {
+ if (dbMetaData.storesLowerCaseIdentifiers()) {
+ return StringUtils.toLowerCase(identifier);
+ }
+ if (dbMetaData.storesUpperCaseIdentifiers()) {
+ return StringUtils.toUpperCase(identifier);
+ }
+ return identifier;
+ }
+
+ /**
+ * Builds the select list for an export query, quoting each column with the
database's
+ * identifier quote string so that reserved words such as "end" are handled
correctly.
+ */
+ protected String makeExportSelectList(final Connection conn, final
List<String> columns) throws SQLException
+ {
+ final String quote = conn.getMetaData().getIdentifierQuoteString();
+ if (quote == null || " ".equals(quote)) {
+ return String.join(",", columns);
+ }
+ return columns.stream()
+ .map(column -> quote + StringUtils.replace(column, quote,
quote + quote) + quote)
+ .collect(Collectors.joining(","));
+ }
+
+ /**
+ * Exports a table to a CSV file using generic JDBC.
+ * Binary columns are hex-encoded and booleans are written as true/false
strings.
+ * Subclasses may override {@link #exportTable} with a database-specific
implementation
+ * while this method remains available for testing or fallback.
+ *
+ * @param columns columns to export in the desired order, or null to export
all columns
+ */
+ protected void exportTableWithJdbc(
+ final String tableName,
+ final String outputPath,
+ @Nullable final List<String> columns
+ )
+ {
+ // Use a transaction so that the connection has autoCommit=false.
+ // PostgreSQL JDBC requires autoCommit=false and a positive fetch size
+ // to use cursor-based streaming instead of buffering the entire ResultSet.
+ retryTransaction(
+ (TransactionCallback<Void>) (handle, status) -> {
+ final Connection conn = handle.getConnection();
+ final String selectList = columns == null || columns.isEmpty() ? "*"
: makeExportSelectList(conn, columns);
+ try (Statement stmt = conn.createStatement()) {
+ final int fetchSize = getStreamingFetchSize();
+ if (fetchSize > 0) {
+ stmt.setFetchSize(fetchSize);
+ }
+ try (ResultSet rs = stmt.executeQuery(StringUtils.format("SELECT
%s FROM %s", selectList, tableName));
+ OutputStreamWriter writer =
+ new OutputStreamWriter(new FileOutputStream(outputPath),
StandardCharsets.UTF_8)) {
+ final int columnCount = rs.getMetaData().getColumnCount();
+ final List<String> values = new ArrayList<>(columnCount);
+ while (rs.next()) {
+ values.clear();
+ for (int i = 1; i <= columnCount; i++) {
+ values.add(readCsvValue(rs, i));
+ }
+ writer.write(String.join(",", values));
+ writer.write('\n');
+ }
+ }
+ }
+ return null;
+ },
+ QUIET_RETRIES,
+ DEFAULT_MAX_TRIES
+ );
+ }
+
+ /**
+ * Reads the given column of the current row as a CSV field. Binary values
are hex-encoded, booleans are written
+ * as true/false and NULLs as empty fields.
+ */
+ private static String readCsvValue(final ResultSet rs, final int column)
throws SQLException
+ {
+ final ResultSetMetaData meta = rs.getMetaData();
+ final int type = meta.getColumnType(column);
+ if (type == Types.BINARY || type == Types.VARBINARY || type ==
Types.LONGVARBINARY || type == Types.BLOB
+ || (type == Types.OTHER &&
"bytea".equalsIgnoreCase(meta.getColumnTypeName(column)))) {
+ final byte[] bytes = rs.getBytes(column);
+ return bytes == null ? "" : BaseEncoding.base16().encode(bytes);
+ } else if (type == Types.BOOLEAN || type == Types.BIT) {
+ final boolean value = rs.getBoolean(column);
+ return rs.wasNull() ? "" : String.valueOf(value);
+ } else {
+ return csvEscape(rs.getString(column));
+ }
+ }
+
+ /**
+ * Escapes a value for CSV output as per RFC 4180: values containing a
comma, double quote or line break are
+ * wrapped in double quotes, with inner double quotes doubled. A null value
is written as an empty field.
+ */
+ public static String csvEscape(@Nullable final String value)
Review Comment:
[P2] MySQL import consumes backslashes
The RFC CSV writer leaves backslashes literal, while the documented MySQL
LOAD DATA command uses MySQL's default backslash escape processing. Values such
as JSON backslash-n or segment identifiers containing backslashes are therefore
altered during PostgreSQL-to-MySQL migration. Use matching escape semantics in
the writer and import command.
--
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]