From 008853e6b9fa63a91cf1d7ffaa15f5280ecafe08 Mon Sep 17 00:00:00 2001 From: Stefan Bischof Date: Sat, 15 Aug 2026 11:14:32 +0200 Subject: [PATCH] feat(dialect,jdbc): bulk-load CSV imports through the dialect Signed-off-by: Stefan Bischof --- .../daanse/sql/dialect/api/Dialect.java | 9 +++ .../api/generator/BulkLoadGenerator.java | 49 ++++++++++++ .../sql/dialect/db/duckdb/DuckDbDialect.java | 30 ++++++++ .../db/postgresql/PostgreSqlDialect.java | 1 + .../sql/dialect/db/sqlite/SqliteDialect.java | 1 - .../importer/csv/impl/CsvDataImporter.java | 11 ++- .../importer/csv/impl/DialectAwareLoader.java | 76 +++++++++++++++++++ 7 files changed, 174 insertions(+), 3 deletions(-) create mode 100644 dialect/api/src/main/java/org/eclipse/daanse/sql/dialect/api/generator/BulkLoadGenerator.java create mode 100644 jdbc/importer/csv/src/main/java/org/eclipse/daanse/sql/jdbc/importer/csv/impl/DialectAwareLoader.java diff --git a/dialect/api/src/main/java/org/eclipse/daanse/sql/dialect/api/Dialect.java b/dialect/api/src/main/java/org/eclipse/daanse/sql/dialect/api/Dialect.java index 7811c3e..78cc5a5 100644 --- a/dialect/api/src/main/java/org/eclipse/daanse/sql/dialect/api/Dialect.java +++ b/dialect/api/src/main/java/org/eclipse/daanse/sql/dialect/api/Dialect.java @@ -157,6 +157,15 @@ default MergeGenerator mergeGenerator() { }; } + /** + * Bulk-load statement emission ({@code CSVREAD}, {@code COPY}, + * {@code LOAD DATA ...}). + */ + default org.eclipse.daanse.sql.dialect.api.generator.BulkLoadGenerator bulkLoadGenerator() { + return new org.eclipse.daanse.sql.dialect.api.generator.BulkLoadGenerator() { + }; + } + /** * Type-cast emission ({@code CAST(x AS T)}, {@code TRY_CAST}, * {@code SAFE_CAST}). diff --git a/dialect/api/src/main/java/org/eclipse/daanse/sql/dialect/api/generator/BulkLoadGenerator.java b/dialect/api/src/main/java/org/eclipse/daanse/sql/dialect/api/generator/BulkLoadGenerator.java new file mode 100644 index 0000000..607f8ed --- /dev/null +++ b/dialect/api/src/main/java/org/eclipse/daanse/sql/dialect/api/generator/BulkLoadGenerator.java @@ -0,0 +1,49 @@ +/* + * Copyright (c) 2026 Contributors to the Eclipse Foundation. + * + * This program and the accompanying materials are made + * available under the terms of the Eclipse Public License 2.0 + * which is available at https://www.eclipse.org/legal/epl-2.0/ + * + * SPDX-License-Identifier: EPL-2.0 + * + * Contributors: + * SmartCity Jena - initial + * Stefan Bischof (bipolis.org) - initial + */ +package org.eclipse.daanse.sql.dialect.api.generator; + +import java.nio.file.Path; +import java.util.List; +import java.util.Optional; + +import org.eclipse.daanse.sql.model.schema.TableReference; + +/** + * Generates the database-native bulk-load statement for a local delimited file + * (H2 {@code CSVREAD}, DuckDB {@code read_csv}, PostgreSQL {@code COPY}, MySQL + * {@code LOAD DATA LOCAL INFILE}). The default is empty; callers fall back to + * batched INSERTs. The file path must be visible to whatever executes the + * statement. + */ +public interface BulkLoadGenerator { + + /** + * The bulk-load statement for a delimited text file whose data begins after + * {@code skipLines} lines. Values are read by position; the caller has already + * created the table with {@code columns} in this order. + * + * @param skipLines lines before the first data line, headers included + * @param nullLiteral text that stands for an absent value + * @return the executable SQL statement, or empty when unsupported + */ + default Optional loadFromDelimitedFile(TableReference target, List columns, Path csvFile, + char delimiter, int skipLines, String nullLiteral) { + return Optional.empty(); + } + + /** Whether this dialect generates native bulk-load statements. */ + default boolean supportsBulkLoad() { + return false; + } +} diff --git a/dialect/db/duckdb/src/main/java/org/eclipse/daanse/sql/dialect/db/duckdb/DuckDbDialect.java b/dialect/db/duckdb/src/main/java/org/eclipse/daanse/sql/dialect/db/duckdb/DuckDbDialect.java index 14b0c0f..8d60bf5 100644 --- a/dialect/db/duckdb/src/main/java/org/eclipse/daanse/sql/dialect/db/duckdb/DuckDbDialect.java +++ b/dialect/db/duckdb/src/main/java/org/eclipse/daanse/sql/dialect/db/duckdb/DuckDbDialect.java @@ -218,4 +218,34 @@ public String paginate(java.util.OptionalLong limit, java.util.OptionalLong offs return local; } + + /** + * Bulk load via {@code read_csv}. Uses {@code header=false} with {@code skip} + * because {@code header=true} would take the type line as the first data row + * and turn every column into VARCHAR. + */ + @Override + public org.eclipse.daanse.sql.dialect.api.generator.BulkLoadGenerator bulkLoadGenerator() { + return new org.eclipse.daanse.sql.dialect.api.generator.BulkLoadGenerator() { + + @Override + public boolean supportsBulkLoad() { + return true; + } + + @Override + public java.util.Optional loadFromDelimitedFile( + org.eclipse.daanse.sql.model.schema.TableReference target, java.util.List columns, + java.nio.file.Path csvFile, char delimiter, int skipLines, String nullLiteral) { + String quotedColumns = columns.stream().map(DuckDbDialect.this::quoteIdentifier) + .collect(java.util.stream.Collectors.joining(", ")); + String file = csvFile.toAbsolutePath().toString().replace("'", "''"); + return java.util.Optional.of("INSERT INTO " + qualified(target) + " (" + quotedColumns + + ") SELECT * FROM read_csv('" + file + "', header=false, skip=" + skipLines + ", delim='" + + (delimiter == '\'' ? "''" : String.valueOf(delimiter)) + "', nullstr='" + + nullLiteral.replace("'", "''") + "')"); + } + }; + } + } diff --git a/dialect/db/postgresql/src/main/java/org/eclipse/daanse/sql/dialect/db/postgresql/PostgreSqlDialect.java b/dialect/db/postgresql/src/main/java/org/eclipse/daanse/sql/dialect/db/postgresql/PostgreSqlDialect.java index 3a0c60f..5adc3cf 100644 --- a/dialect/db/postgresql/src/main/java/org/eclipse/daanse/sql/dialect/db/postgresql/PostgreSqlDialect.java +++ b/dialect/db/postgresql/src/main/java/org/eclipse/daanse/sql/dialect/db/postgresql/PostgreSqlDialect.java @@ -397,4 +397,5 @@ public boolean supportsPercentileCont() { public boolean supportsNthValue() { return true; } + } diff --git a/dialect/db/sqlite/src/main/java/org/eclipse/daanse/sql/dialect/db/sqlite/SqliteDialect.java b/dialect/db/sqlite/src/main/java/org/eclipse/daanse/sql/dialect/db/sqlite/SqliteDialect.java index 648d119..71b6aff 100644 --- a/dialect/db/sqlite/src/main/java/org/eclipse/daanse/sql/dialect/db/sqlite/SqliteDialect.java +++ b/dialect/db/sqlite/src/main/java/org/eclipse/daanse/sql/dialect/db/sqlite/SqliteDialect.java @@ -265,5 +265,4 @@ public java.util.Optional upsert(UpsertSpec spec, java.util.List cachedMergeGenerator = local; return local; } - } diff --git a/jdbc/importer/csv/src/main/java/org/eclipse/daanse/sql/jdbc/importer/csv/impl/CsvDataImporter.java b/jdbc/importer/csv/src/main/java/org/eclipse/daanse/sql/jdbc/importer/csv/impl/CsvDataImporter.java index fc52364..de5a8cd 100644 --- a/jdbc/importer/csv/src/main/java/org/eclipse/daanse/sql/jdbc/importer/csv/impl/CsvDataImporter.java +++ b/jdbc/importer/csv/src/main/java/org/eclipse/daanse/sql/jdbc/importer/csv/impl/CsvDataImporter.java @@ -199,7 +199,11 @@ private void loadTable(Connection connection, Path path) throws SQLException { List headersTypeList = getHeadersTypeList(types); if (it.hasNext()) { createTable(connection, headersTypeList, tableDefinition); - insertTable(connection, it, headersTypeList, tableRef); + // skipLines = 2: the column-name line and the SQL-type line. + if (!DialectAwareLoader.loadNatively(connection, dialect, tableRef, headersTypeList, path, + config.fieldSeparator(), 2, config.nullValue())) { + insertTable(connection, it, headersTypeList, tableRef); + } } } catch (IOException e) { @@ -241,7 +245,10 @@ public void createTable(Connection connection, List headersTyp LOGGER.debug("Created table in given database. {}", sql); stmt.execute(sql); - connection.commit(); + if (!connection.getAutoCommit()) { + // commit() on an auto-commit connection is a JDBC error; DuckDB rejects it. + connection.commit(); + } } catch (SQLException e) { throw new CsvDataImporterException("Exception wile create table", e); } diff --git a/jdbc/importer/csv/src/main/java/org/eclipse/daanse/sql/jdbc/importer/csv/impl/DialectAwareLoader.java b/jdbc/importer/csv/src/main/java/org/eclipse/daanse/sql/jdbc/importer/csv/impl/DialectAwareLoader.java new file mode 100644 index 0000000..eb6a642 --- /dev/null +++ b/jdbc/importer/csv/src/main/java/org/eclipse/daanse/sql/jdbc/importer/csv/impl/DialectAwareLoader.java @@ -0,0 +1,76 @@ +/* + * Copyright (c) 2026 Contributors to the Eclipse Foundation. + * + * This program and the accompanying materials are made + * available under the terms of the Eclipse Public License 2.0 + * which is available at https://www.eclipse.org/legal/epl-2.0/ + * + * SPDX-License-Identifier: EPL-2.0 + * + * Contributors: + * SmartCity Jena - initial + * Stefan Bischof (bipolis.org) - initial + */ +package org.eclipse.daanse.sql.jdbc.importer.csv.impl; + +import java.nio.file.Path; +import java.sql.Connection; +import java.sql.SQLException; +import java.sql.Statement; +import java.util.List; +import java.util.Optional; + +import org.eclipse.daanse.sql.dialect.api.Dialect; +import org.eclipse.daanse.sql.model.schema.ColumnDefinition; +import org.eclipse.daanse.sql.model.schema.TableReference; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * Loads a delimited file with the dialect's native bulk-load statement, if it + * has one; otherwise the caller falls back to batched INSERTs. + */ +final class DialectAwareLoader { + + private static final Logger LOGGER = LoggerFactory.getLogger(DialectAwareLoader.class); + + private DialectAwareLoader() { + // static access only + } + + /** + * @param skipLines lines before the first data line + * @return whether the file was loaded; {@code false} means the caller has to + * do it row by row + */ + static boolean loadNatively(Connection connection, Dialect dialect, TableReference target, + List columns, Path file, char delimiter, int skipLines, String nullLiteral) + throws SQLException { + if (!dialect.bulkLoadGenerator().supportsBulkLoad()) { + return false; + } + Optional statement = dialect.bulkLoadGenerator().loadFromDelimitedFile(target, + columns.stream().map(column -> column.column().name()).toList(), file, delimiter, skipLines, nullLiteral); + if (statement.isEmpty()) { + LOGGER.debug("{} has no bulk load for a file with {} leading lines; loading {} row by row", + dialect.name(), skipLines, file.getFileName()); + return false; + } + long started = System.currentTimeMillis(); + try (Statement direct = connection.createStatement()) { + direct.execute(statement.get()); + } catch (SQLException e) { + // Typically the server cannot see the file or local reads are off. + LOGGER.warn("{} refused to read {} itself ({}); loading row by row", dialect.name(), file.getFileName(), + e.getMessage()); + // Truncate whatever the failed attempt left, or the fallback load adds up. + try (Statement direct = connection.createStatement()) { + direct.executeUpdate(dialect.ddlGenerator().truncate(target)); + } + return false; + } + LOGGER.info("{} read {} itself in {} ms", dialect.name(), file.getFileName(), + System.currentTimeMillis() - started); + return true; + } +}