From 22814b4a6e9569c0ef8092a8095c552e724fca84 Mon Sep 17 00:00:00 2001 From: Arun Sharma Date: Fri, 2 Oct 2026 22:24:48 -0700 Subject: [PATCH] adbc: add CATALOG support with qualified scans and schema probe fallback - Accept optional CATALOG attach option for 3-level namespaces (e.g. Databricks Unity Catalog), threaded through ADBCCatalog and ADBCConnector::qualifiedTableRef. - Build a prequalified FROM clause once per table and propagate it via ADBCTableScanInfo (fromClause + backtickIds) through catalog entry and scan bind, avoiding double-quoting. - Use backtick quoting for catalog-qualified refs and columns: with ANSI mode off (Spark/Databricks), double-quoted ids are string literals and fail with PARSE_SYNTAX_ERROR. - Fall back to SELECT * WHERE 1=0 statement probe when AdbcConnectionGetTableSchema is unsupported (e.g. Databricks), reading the Arrow stream schema instead. - Switch connector mutex to recursive_mutex since schema probe re-enters via executeQuery under the same lock. --- adbc/src/catalog/adbc_catalog.cpp | 5 +- adbc/src/catalog/adbc_table_catalog_entry.cpp | 3 +- adbc/src/connector/adbc_connector.cpp | 105 +++++++++++++++--- adbc/src/function/adbc_scan.cpp | 23 +++- adbc/src/include/catalog/adbc_catalog.h | 6 +- adbc/src/include/connector/adbc_connector.h | 8 +- adbc/src/include/function/adbc_scan.h | 16 ++- adbc/src/storage/adbc_storage.cpp | 14 ++- 8 files changed, 152 insertions(+), 28 deletions(-) diff --git a/adbc/src/catalog/adbc_catalog.cpp b/adbc/src/catalog/adbc_catalog.cpp index 47955601..97a13880 100644 --- a/adbc/src/catalog/adbc_catalog.cpp +++ b/adbc/src/catalog/adbc_catalog.cpp @@ -25,8 +25,9 @@ void ADBCCatalog::createForeignTable(const std::string& tableName) { columnNames.push_back(name); columnTypes.push_back(type.copy()); } - auto scanInfo = std::make_shared(tableName, columnNames, - copyVector(columnTypes), connector); + auto scanInfo = std::make_shared(tableName, + connector.qualifiedTableRef(catalogName, schemaName, tableName), + !catalogName.empty() /*backtickIds*/, columnNames, copyVector(columnTypes), connector); auto attachedEntry = std::make_unique(tableName, getADBCScanFunction(scanInfo), scanInfo); for (auto i = 0u; i < columnNames.size(); i++) { diff --git a/adbc/src/catalog/adbc_table_catalog_entry.cpp b/adbc/src/catalog/adbc_table_catalog_entry.cpp index 31200661..254b9a6f 100644 --- a/adbc/src/catalog/adbc_table_catalog_entry.cpp +++ b/adbc/src/catalog/adbc_table_catalog_entry.cpp @@ -38,7 +38,8 @@ std::unique_ptr ADBCTableCatalogEntry::getBoundScanI scanColumnTypes.push_back(scanInfo->columnTypes[i].copy()); } auto boundScanInfo = std::make_shared(scanInfo->tableName, - std::move(scanColumnNames), std::move(scanColumnTypes), scanInfo->connector); + scanInfo->fromClause, scanInfo->backtickIds, std::move(scanColumnNames), + std::move(scanColumnTypes), scanInfo->connector); auto bindData = std::make_unique(std::move(boundScanInfo), std::move(columns)); return std::make_unique(scanFunction, std::move(bindData)); diff --git a/adbc/src/connector/adbc_connector.cpp b/adbc/src/connector/adbc_connector.cpp index ac4125b1..1658f5d1 100644 --- a/adbc/src/connector/adbc_connector.cpp +++ b/adbc/src/connector/adbc_connector.cpp @@ -13,6 +13,39 @@ namespace adbc_extension { static constexpr const char* DRIVER_OPTION = "DRIVER"; static constexpr const char* TABLES_OPTION = "TABLES"; static constexpr const char* SCHEMA_OPTION = "SCHEMA"; +static constexpr const char* CATALOG_OPTION = "CATALOG"; + +// Backtick-quote one SQL identifier part (internal backticks are doubled). +// Parts are quoted separately so catalog.schema.table keeps its structure. +// Backticks (not double quotes) because the qualified path targets lakehouse +// engines (Spark/Databricks): with ANSI mode off, "x" is a string literal +// and SELECT * FROM "catalog"."schema"."table" fails with +// PARSE_SYNTAX_ERROR. DuckDB accepts backticks as well. +static std::string quoteSQLPart(const std::string& part) { + std::string out = "`"; + for (const char c : part) { + out += c; + if (c == '`') { + out += c; + } + } + out += '`'; + return out; +} + +// Legacy double-quote for the unqualified path (byte-identical to the old +// quoteIdentifier behavior in adbc_scan.cpp). +static std::string quoteLegacy(const std::string& part) { + std::string out = "\""; + for (const char c : part) { + out += c; + if (c == '"') { + out += c; + } + } + out += "\""; + return out; +} static std::vector splitCommaSeparated(const std::string& input) { std::vector result; @@ -116,7 +149,8 @@ void ADBCConnector::connect(const std::string& uri) { } for (auto& [key, value] : attachOption.options) { auto upperKey = common::StringUtils::getUpper(key); - if (upperKey == DRIVER_OPTION || upperKey == TABLES_OPTION || upperKey == SCHEMA_OPTION) { + if (upperKey == DRIVER_OPTION || upperKey == TABLES_OPTION || upperKey == SCHEMA_OPTION || + upperKey == CATALOG_OPTION) { continue; } if (value.getDataType().getLogicalTypeID() != common::LogicalTypeID::STRING) { @@ -141,9 +175,34 @@ std::vector ADBCConnector::getTableNames() const { return splitCommaSeparated(tables); } +std::string ADBCConnector::getCatalogName() const { + return getStringOption(CATALOG_OPTION); +} + +std::string ADBCConnector::qualifiedTableRef(const std::string& catalog, const std::string& schema, + const std::string& table) const { + if (catalog.empty()) { + // Legacy behavior: bare quoted table (schema resolved by the driver). + return quoteLegacy(table); + } + std::string ref; + bool first = true; + for (const auto* part : {&catalog, &schema, &table}) { + if (part->empty()) { + continue; + } + if (!first) { + ref += "."; + } + ref += quoteSQLPart(*part); + first = false; + } + return ref.empty() ? quoteSQLPart(table) : ref; +} + std::vector> ADBCConnector::getTableSchema( const std::string& schemaName, const std::string& tableName) const { - std::lock_guard lock{mtx}; + std::lock_guard lock{mtx}; if (schemaCache.contains(tableName)) { std::vector> cachedResult; cachedResult.reserve(schemaCache.at(tableName).size()); @@ -152,17 +211,37 @@ std::vector> ADBCConnector::getTable } return cachedResult; } - ArrowSchemaWrapper schema; - checkStatus(AdbcConnectionGetTableSchema(&connection, nullptr, - schemaName.empty() ? nullptr : schemaName.c_str(), tableName.c_str(), &schema, - &error), - std::format("AdbcConnectionGetTableSchema({})", tableName)); std::vector> result; - result.reserve(schema.n_children); - for (auto i = 0; i < schema.n_children; i++) { - auto child = schema.children[i]; - result.emplace_back(child->name == nullptr ? std::format("column{}", i) : child->name, - common::ArrowConverter::fromArrowSchema(child)); + bool haveSchema = false; + try { + ArrowSchemaWrapper schema; + checkStatus(AdbcConnectionGetTableSchema(&connection, nullptr, + schemaName.empty() ? nullptr : schemaName.c_str(), tableName.c_str(), + &schema, &error), + std::format("AdbcConnectionGetTableSchema({})", tableName)); + result.reserve(schema.n_children); + for (auto i = 0; i < schema.n_children; i++) { + auto child = schema.children[i]; + result.emplace_back(child->name == nullptr ? std::format("column{}", i) : child->name, + common::ArrowConverter::fromArrowSchema(child)); + } + haveSchema = true; + } catch (const common::Exception&) { + // Drivers without AdbcConnectionGetTableSchema (e.g. Databricks) + // fall through to the statement probe below. + } + if (!haveSchema) { + // Schema-on-first-scan: run an empty result query and read the + // Arrow stream schema. Needs no driver metadata API. + auto probe = executeQuery(std::format("SELECT * FROM {} WHERE 1=0", + qualifiedTableRef(getCatalogName(), schemaName, tableName)), + {} /*columnNames*/, {} /*columnTypes*/); + result.reserve(probe->schema.n_children); + for (auto i = 0; i < probe->schema.n_children; i++) { + auto child = probe->schema.children[i]; + result.emplace_back(child->name == nullptr ? std::format("column{}", i) : child->name, + common::ArrowConverter::fromArrowSchema(child)); + } } std::vector> cachedResult; cachedResult.reserve(result.size()); @@ -190,7 +269,7 @@ std::unique_ptr ADBCConnector::executeQuery(const std::string& throw common::RuntimeException{std::format("{} failed: {}", operation, message)}; }; { - std::lock_guard lock{mtx}; + std::lock_guard lock{mtx}; checkQueryStatus(AdbcConnectionNew(&result->connection, &queryError), "AdbcConnectionNew"); result->connectionInitialized = true; checkQueryStatus(AdbcConnectionInit(&result->connection, diff --git a/adbc/src/function/adbc_scan.cpp b/adbc/src/function/adbc_scan.cpp index 9739a703..c21fff45 100644 --- a/adbc/src/function/adbc_scan.cpp +++ b/adbc/src/function/adbc_scan.cpp @@ -26,22 +26,34 @@ static std::string quoteIdentifier(const std::string& value) { return result; } -static std::string joinColumns(const std::vector& columnNames) { +static std::string quoteBacktick(const std::string& value) { + std::string result = "`"; + for (auto ch : value) { + result += ch; + if (ch == '`') { + result += ch; + } + } + result += "`"; + return result; +} + +static std::string joinColumns(const std::vector& columnNames, bool backtickIds) { std::string result; bool first = true; for (auto& columnName : columnNames) { if (!first) { result += ", "; } - result += quoteIdentifier(columnName); + result += backtickIds ? quoteBacktick(columnName) : quoteIdentifier(columnName); first = false; } return result.empty() ? "*" : result; } std::string ADBCScanBindData::getSQL() const { - auto sql = std::format("SELECT {} FROM {}", joinColumns(scanInfo->columnNames), - quoteIdentifier(scanInfo->tableName)); + auto sql = std::format("SELECT {} FROM {}", + joinColumns(scanInfo->columnNames, scanInfo->backtickIds), scanInfo->fromClause); if (getLimitNum() != common::INVALID_ROW_IDX) { sql += std::format(" LIMIT {}", getLimitNum()); } @@ -126,7 +138,8 @@ std::unique_ptr ADBCScanFunction::bindFunc( } auto columns = input->binder->createVariables(columnNames, columnTypes); auto selectedScanInfo = std::make_shared(scanInfo->tableName, - std::move(columnNames), std::move(columnTypes), scanInfo->connector); + scanInfo->fromClause, scanInfo->backtickIds, std::move(columnNames), std::move(columnTypes), + scanInfo->connector); return std::make_unique(std::move(selectedScanInfo), std::move(columns)); } diff --git a/adbc/src/include/catalog/adbc_catalog.h b/adbc/src/include/catalog/adbc_catalog.h index 2402cb86..4605826f 100644 --- a/adbc/src/include/catalog/adbc_catalog.h +++ b/adbc/src/include/catalog/adbc_catalog.h @@ -9,9 +9,10 @@ namespace adbc_extension { class ADBCCatalog final : public extension::CatalogExtension { public: - ADBCCatalog(std::string schemaName, main::ClientContext* context, + ADBCCatalog(std::string catalogName, std::string schemaName, main::ClientContext* context, const ADBCConnector& connector) - : schemaName{std::move(schemaName)}, context{context}, connector{connector} {} + : catalogName{std::move(catalogName)}, schemaName{std::move(schemaName)}, context{context}, + connector{connector} {} void init() override; @@ -19,6 +20,7 @@ class ADBCCatalog final : public extension::CatalogExtension { void createForeignTable(const std::string& tableName); private: + std::string catalogName; std::string schemaName; main::ClientContext* context; const ADBCConnector& connector; diff --git a/adbc/src/include/connector/adbc_connector.h b/adbc/src/include/connector/adbc_connector.h index f312a116..8d7217d3 100644 --- a/adbc/src/include/connector/adbc_connector.h +++ b/adbc/src/include/connector/adbc_connector.h @@ -1,6 +1,7 @@ #pragma once #include +#include #include #include "binder/bound_attach_info.h" @@ -55,6 +56,9 @@ class ADBCConnector final { std::vector getTableNames() const; std::vector> getTableSchema( const std::string& schemaName, const std::string& tableName) const; + std::string getCatalogName() const; + std::string qualifiedTableRef(const std::string& catalog, const std::string& schema, + const std::string& table) const; std::unique_ptr executeQuery(const std::string& query, const std::vector& columnNames, const std::vector& columnTypes) const; @@ -66,7 +70,9 @@ class ADBCConnector final { private: const binder::AttachOption& attachOption; - mutable std::mutex mtx; + // Recursive: getTableSchema() may fall back to a statement probe that + // re-enters through the same locking as executeQuery(). + mutable std::recursive_mutex mtx; mutable AdbcError error{}; AdbcDatabase database{}; mutable AdbcConnection connection{}; diff --git a/adbc/src/include/function/adbc_scan.h b/adbc/src/include/function/adbc_scan.h index 6a9545ae..c03aac68 100644 --- a/adbc/src/include/function/adbc_scan.h +++ b/adbc/src/include/function/adbc_scan.h @@ -10,13 +10,23 @@ namespace adbc_extension { struct ADBCTableScanInfo { std::string tableName; + // Prebuilt FROM clause (catalog/schema-qualified when the ATTACH + // specified CATALOG, else the bare quoted table). Built once so the + // quoting is not applied twice. + std::string fromClause; + // Lakehouse-quoting for identifiers: with ANSI mode off (Spark/ + // Databricks), "x" is a string literal, so columns must use backticks. + // Legacy drivers keep double quotes. + bool backtickIds = false; std::vector columnNames; std::vector columnTypes; const ADBCConnector& connector; - ADBCTableScanInfo(std::string tableName, std::vector columnNames, - std::vector columnTypes, const ADBCConnector& connector) - : tableName{std::move(tableName)}, columnNames{std::move(columnNames)}, + ADBCTableScanInfo(std::string tableName, std::string fromClause, bool backtickIds, + std::vector columnNames, std::vector columnTypes, + const ADBCConnector& connector) + : tableName{std::move(tableName)}, fromClause{std::move(fromClause)}, + backtickIds{backtickIds}, columnNames{std::move(columnNames)}, columnTypes{std::move(columnTypes)}, connector{connector} {} }; diff --git a/adbc/src/storage/adbc_storage.cpp b/adbc/src/storage/adbc_storage.cpp index ecbf9ec5..3c4db37b 100644 --- a/adbc/src/storage/adbc_storage.cpp +++ b/adbc/src/storage/adbc_storage.cpp @@ -23,9 +23,21 @@ std::unique_ptr attachADBC(std::string dbName, std::stri } schemaName = val.getValue(); } + // Optional catalog for 3-level namespaces (e.g. Databricks UC). Enables + // the statement-probe schema discovery and qualifies scan SQL. Empty = + // legacy behavior (bare table, driver-resolved namespace). + std::string catalogName; + if (attachOption.options.contains("CATALOG")) { + auto val = attachOption.options.at("CATALOG"); + if (val.getDataType().getLogicalTypeID() != common::LogicalTypeID::STRING) { + throw common::RuntimeException{"Invalid option value for CATALOG"}; + } + catalogName = val.getValue(); + } auto connector = std::make_unique(attachOption); connector->connect(dbPath); - auto catalog = std::make_unique(schemaName, clientContext, *connector); + auto catalog = + std::make_unique(catalogName, schemaName, clientContext, *connector); catalog->init(); return std::make_unique(dbName, ADBCStorageExtension::DB_TYPE, std::move(catalog), std::move(connector));