Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
48 changes: 39 additions & 9 deletions duckdb/src/catalog/duckdb_catalog.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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<BoundExtraCreateDuckDBTableInfo>();
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) {
Expand Down Expand Up @@ -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;
Expand All @@ -202,7 +204,9 @@ void DuckDBCatalog::createForeignRelTable(const std::string& tableName, bool int

// Build property definitions
std::vector<binder::PropertyDefinition> 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
Expand Down Expand Up @@ -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<DuckDBTableScanInfo>(queryStr, std::move(columnTypes),
columnNames, connector);
auto scanFunc = getScanFunction(duckdbTableInfo);
Expand Down Expand Up @@ -280,7 +284,6 @@ void DuckDBCatalog::createForeignRelTable(const std::string& tableName, bool int
auto foreignDatabaseName = dbName;

std::vector<catalog::RelTableCatalogInfo> 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);
Expand Down Expand Up @@ -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<duckdb::MaterializedQueryResult> 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++) {
Expand Down
53 changes: 53 additions & 0 deletions duckdb/src/include/catalog/duckdb_schema_utils.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,53 @@
#pragma once

#include <algorithm>
#include <iterator>
#include <string>

#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<std::string>() == "__" &&
result.GetValue(std::distance(result.names.begin(), typeColumn), 0)
.GetValue<std::string>() == "UNKNOWN";
}

} // namespace duckdb_extension
} // namespace lbug
40 changes: 25 additions & 15 deletions duckdb/src/include/storage/attached_duckdb_database.h
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@
#include <string>
#include <vector>

#include "catalog/duckdb_schema_utils.h"
#include "connector/duckdb_connector.h"
#include "main/attached_database.h"
#include <format>
Expand All @@ -22,8 +23,14 @@ inline std::vector<std::string> splitQualifiedTableName(const std::string& table
std::vector<std::string> 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;
}
Expand All @@ -38,19 +45,6 @@ inline std::vector<std::string> 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,
Expand Down Expand Up @@ -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<std::string>()),
quoteDuckDBIdentifier(result->GetValue(3, 0).GetValue<std::string>()),
quoteDuckDBIdentifier(unqualified));
try {
return connector
->executeQuery(std::format("SELECT * FROM {} LIMIT 0", qualifiedName))
->names;
} catch (const common::Exception&) {
return {};
}
}
std::vector<std::string> columnNames;
for (auto i = 0u; i < result->RowCount(); i++) {
columnNames.push_back(result->GetValue(0, i).GetValue<std::string>());
Expand Down
9 changes: 6 additions & 3 deletions iceberg/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
Loading