From bb7186563b65e9d0596c087e743d0a926f3c377f Mon Sep 17 00:00:00 2001 From: ericyuanhui <285521263@qq.com> Date: Wed, 23 Sep 2026 22:00:32 +0800 Subject: [PATCH] fix attaching an Iceberg database through a REST catalog Signed-off-by: ericyuanhui <285521263@qq.com> --- duckdb/src/catalog/duckdb_catalog.cpp | 48 +++++++++++++---- .../src/include/catalog/duckdb_schema_utils.h | 53 +++++++++++++++++++ .../storage/attached_duckdb_database.h | 40 ++++++++------ iceberg/README.md | 9 ++-- 4 files changed, 123 insertions(+), 27 deletions(-) create mode 100644 duckdb/src/include/catalog/duckdb_schema_utils.h diff --git a/duckdb/src/catalog/duckdb_catalog.cpp b/duckdb/src/catalog/duckdb_catalog.cpp index f3fdca6..1eb7e91 100644 --- a/duckdb/src/catalog/duckdb_catalog.cpp +++ b/duckdb/src/catalog/duckdb_catalog.cpp @@ -8,6 +8,7 @@ #include "binder/expression/variable_expression.h" #include "catalog/catalog_entry/node_table_catalog_entry.h" #include "catalog/catalog_entry/rel_group_catalog_entry.h" +#include "catalog/duckdb_schema_utils.h" #include "catalog/duckdb_table_catalog_entry.h" #include "common/exception/binder.h" #include "common/exception/runtime.h" @@ -43,7 +44,7 @@ void DuckDBCatalog::init() { auto query = std::format( "select table_name from information_schema.tables where table_catalog = '{}' and " "table_schema = '{}' order by table_name;", - catalogName, defaultSchemaName); + escapeSingleQuotes(catalogName), escapeSingleQuotes(defaultSchemaName)); auto result = connector.executeQuery(query); // Collect every table name across all result chunks: catalogs with more // tables than fit a single DuckDB data chunk must not lose tables. @@ -113,8 +114,9 @@ std::string DuckDBCatalog::bindSchemaName(const binder::AttachOption& options, static std::string getQuery(const binder::BoundCreateTableInfo& info) { auto extraInfo = info.extraInfo->constPtrCast(); - return "SELECT {} " + std::format("FROM \"{}\".{}.{}", extraInfo->catalogName, - extraInfo->schemaName, info.tableName); + return "SELECT {} " + + std::format("FROM {}.{}.{}", quoteDuckDBIdentifier(extraInfo->catalogName), + quoteDuckDBIdentifier(extraInfo->schemaName), quoteDuckDBIdentifier(info.tableName)); } void DuckDBCatalog::createForeignTable(const std::string& tableName) { @@ -175,7 +177,7 @@ void DuckDBCatalog::createForeignRelTable(const std::string& tableName, bool int "referenced_table FROM duckdb_constraints() " "WHERE constraint_type = 'FOREIGN KEY' " "AND referenced_table IS NOT NULL AND table_name = '{}'", - tableName); + escapeSingleQuotes(tableName)); auto fkResult = connector.executeQuery(fkQuery); std::string srcTableName, dstTableName; @@ -202,7 +204,9 @@ void DuckDBCatalog::createForeignRelTable(const std::string& tableName, bool int // Build property definitions std::vector propertyDefinitions; - bindPropertyDefinitions(tableName, propertyDefinitions); + if (!bindPropertyDefinitions(tableName, propertyDefinitions)) { + return; + } // Determine the node table IDs from the main catalog. containsTable() must // be checked first: getTableCatalogEntry() throws when the table is @@ -238,8 +242,8 @@ void DuckDBCatalog::createForeignRelTable(const std::string& tableName, bool int } // Build query and scan info - auto queryStr = - std::format("SELECT * FROM \"{}\".{}.{}", catalogName, defaultSchemaName, tableName); + auto queryStr = std::format("SELECT * FROM {}.{}.{}", quoteDuckDBIdentifier(catalogName), + quoteDuckDBIdentifier(defaultSchemaName), quoteDuckDBIdentifier(tableName)); auto duckdbTableInfo = std::make_shared(queryStr, std::move(columnTypes), columnNames, connector); auto scanFunc = getScanFunction(duckdbTableInfo); @@ -280,7 +284,6 @@ void DuckDBCatalog::createForeignRelTable(const std::string& tableName, bool int auto foreignDatabaseName = dbName; std::vector relTableInfos; - auto info = bindCreateTableInfo(tableName); common::oid_t relOID = tables->getNextOID(); relTableInfos.emplace_back(catalog::NodeTableIDPair{srcTableID, dstTableID}, relOID, common::RelMultiplicity::MANY, common::RelMultiplicity::MANY); @@ -313,11 +316,38 @@ static bool getTableInfo(const DuckDBConnector& connector, const std::string& ta auto query = std::format("select data_type,column_name from information_schema.columns where " "table_name = '{}' and table_schema = '{}' and table_catalog = '{}' " "order by ordinal_position;", - tableName, schemaName, catalogName); + escapeSingleQuotes(tableName), escapeSingleQuotes(schemaName), + escapeSingleQuotes(catalogName)); auto result = connector.executeQuery(query); if (result->RowCount() == 0) { return false; } + if (isPlaceholderSchema(*result)) { + std::unique_ptr schemaResult; + try { + schemaResult = connector.executeQuery( + std::format("SELECT * FROM {}.{}.{} LIMIT 0", quoteDuckDBIdentifier(catalogName), + quoteDuckDBIdentifier(schemaName), quoteDuckDBIdentifier(tableName))); + } catch (const common::Exception&) { + if (skipUnsupportedTable) { + return false; + } + throw; + } + columnTypes.reserve(schemaResult->types.size()); + columnNames = schemaResult->names; + for (auto& type : schemaResult->types) { + try { + columnTypes.push_back(DuckDBTypeConverter::convertDuckDBType(type.ToString())); + } catch (common::BinderException&) { + if (skipUnsupportedTable) { + return false; + } + throw; + } + } + return true; + } columnTypes.reserve(result->RowCount()); columnNames.reserve(result->RowCount()); for (auto i = 0u; i < result->RowCount(); i++) { diff --git a/duckdb/src/include/catalog/duckdb_schema_utils.h b/duckdb/src/include/catalog/duckdb_schema_utils.h new file mode 100644 index 0000000..cbd7ba1 --- /dev/null +++ b/duckdb/src/include/catalog/duckdb_schema_utils.h @@ -0,0 +1,53 @@ +#pragma once + +#include +#include +#include + +#include "connector/duckdb_connector.h" + +namespace lbug { +namespace duckdb_extension { + +inline std::string escapeSingleQuotes(const std::string& value) { + std::string result; + result.reserve(value.size()); + for (auto c : value) { + if (c == '\'') { + result += "''"; + } else { + result += c; + } + } + return result; +} + +inline std::string quoteDuckDBIdentifier(const std::string& value) { + std::string result = "\""; + for (auto c : value) { + result += c; + if (c == '"') { + result += '"'; + } + } + result += '"'; + return result; +} + +inline bool isPlaceholderSchema(duckdb::MaterializedQueryResult& result) { + if (result.RowCount() != 1) { + return false; + } + const auto nameColumn = std::find(result.names.begin(), result.names.end(), "column_name"); + const auto typeColumn = std::find(result.names.begin(), result.names.end(), "data_type"); + if (nameColumn == result.names.end() || typeColumn == result.names.end()) { + return false; + } + return result.GetValue(std::distance(result.names.begin(), nameColumn), 0) + .GetValue() == "__" && + result.GetValue(std::distance(result.names.begin(), typeColumn), 0) + .GetValue() == "UNKNOWN"; +} + +} // namespace duckdb_extension +} // namespace lbug diff --git a/duckdb/src/include/storage/attached_duckdb_database.h b/duckdb/src/include/storage/attached_duckdb_database.h index 2ffba3e..8192506 100644 --- a/duckdb/src/include/storage/attached_duckdb_database.h +++ b/duckdb/src/include/storage/attached_duckdb_database.h @@ -2,6 +2,7 @@ #include #include +#include "catalog/duckdb_schema_utils.h" #include "connector/duckdb_connector.h" #include "main/attached_database.h" #include @@ -22,8 +23,14 @@ inline std::vector splitQualifiedTableName(const std::string& table std::vector parts; std::string current; bool inQuotes = false; - for (auto c : tableName) { + for (auto i = 0u; i < tableName.size(); i++) { + auto c = tableName[i]; if (c == '"') { + if (inQuotes && i + 1 < tableName.size() && tableName[i + 1] == '"') { + current += '"'; + i++; + continue; + } inQuotes = !inQuotes; continue; } @@ -38,19 +45,6 @@ inline std::vector splitQualifiedTableName(const std::string& table return parts; } -inline std::string escapeSingleQuotes(const std::string& value) { - std::string result; - result.reserve(value.size()); - for (auto c : value) { - if (c == '\'') { - result += "''"; - } else { - result += c; - } - } - return result; -} - class AttachedDuckDBDatabase : public main::AttachedDatabase { public: AttachedDuckDBDatabase(std::string dbName, std::string dbType, @@ -94,11 +88,27 @@ class AttachedDuckDBDatabase : public main::AttachedDatabase { filterCandidates.emplace_back(""); for (auto& filters : filterCandidates) { std::string query = - std::format("SELECT column_name FROM information_schema.columns WHERE table_name " + std::format("SELECT column_name, data_type, table_catalog, table_schema " + "FROM information_schema.columns WHERE table_name " "= '{}'{} ORDER BY ordinal_position", escapeSingleQuotes(unqualified), filters); auto result = connector->executeQuery(query); if (result && result->RowCount() != 0) { + if (isPlaceholderSchema(*result)) { + // Use the catalog and schema of the matched table, including + // for two-part references outside DuckDB's search path. + const auto qualifiedName = std::format("{}.{}.{}", + quoteDuckDBIdentifier(result->GetValue(2, 0).GetValue()), + quoteDuckDBIdentifier(result->GetValue(3, 0).GetValue()), + quoteDuckDBIdentifier(unqualified)); + try { + return connector + ->executeQuery(std::format("SELECT * FROM {} LIMIT 0", qualifiedName)) + ->names; + } catch (const common::Exception&) { + return {}; + } + } std::vector columnNames; for (auto i = 0u; i < result->RowCount(); i++) { columnNames.push_back(result->GetValue(0, i).GetValue()); diff --git a/iceberg/README.md b/iceberg/README.md index 1adcbbc..a7e3ebc 100644 --- a/iceberg/README.md +++ b/iceberg/README.md @@ -45,8 +45,9 @@ is prefixed with the alias automatically). Any other 3-part name is rejected. ## Attaching a REST catalog as a database -An Iceberg REST catalog is a SQL engine, so besides `LOAD FROM` it can be -attached as a database and queried with graph patterns. Attached tables behave +The REST catalog resolves Iceberg metadata; the embedded DuckDB instance executes +SQL against the tables. Besides `LOAD FROM`, the catalog can be attached as a +database and queried with graph patterns. Attached tables behave like any other foreign tables: `MATCH` scans push filters, projections, limits and ordering down to the catalog, and multi-hop patterns over relationship tables are rewritten into a single SQL join (see the foreign join @@ -89,7 +90,9 @@ keeping credentials out of scripts. Authentication is optional: catalogs reachable without credentials (e.g. `s3_tables`/`glue` combined with S3 environment credentials) can be attached -with just `iceberg_warehouse` and `iceberg_endpoint(_type)`. Data files on S3 +without a token. For an unauthenticated REST endpoint, set +`CALL iceberg_authorization_type='none';` before attaching; DuckDB otherwise +defaults to OAuth2 and requires credentials. Data files on S3 are read with the usual `s3_*` options of the httpfs integration. ### Time travel