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
50 changes: 50 additions & 0 deletions src/Common/ProfileEvents.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -1724,6 +1724,56 @@ The server successfully detected this situation and will download merged part fr
M(StatelessWorkerDiscoveryHeartbeatsRejected, "Number of heartbeats the stateless worker discovery service rejected because the worker had already been evicted.", ValueType::Number) \
M(StatelessWorkerDiscoveryKeeperTransactionRetries, "Number of write transactions the stateless worker discovery service retried because its coordination store (Keeper) state was modified concurrently.", ValueType::Number) \
\
M(DataLakeRestCatalogLoadConfig, "Number of 'load config' requests to Iceberg REST catalog.", ValueType::Number) \
M(DataLakeRestCatalogLoadConfigMicroseconds, "Total time of 'load config' requests to Iceberg REST catalog.", ValueType::Microseconds) \
M(DataLakeRestCatalogGetNamespaces, "Number of 'get namespaces' requests to Iceberg REST catalog.", ValueType::Number) \
M(DataLakeRestCatalogGetNamespacesMicroseconds, "Total time of 'get namespaces' requests to Iceberg REST catalog.", ValueType::Microseconds) \
M(DataLakeRestCatalogGetTables, "Number of 'get tables' requests to Iceberg REST catalog.", ValueType::Number) \
M(DataLakeRestCatalogGetTablesMicroseconds, "Total time of 'get tables' requests to Iceberg REST catalog.", ValueType::Microseconds) \
M(DataLakeRestCatalogGetTableMetadata, "Number of 'get table metadata' requests to Iceberg REST catalog.", ValueType::Number) \
M(DataLakeRestCatalogGetTableMetadataMicroseconds, "Total time of 'get table metadata' requests to Iceberg REST catalog.", ValueType::Microseconds) \
M(DataLakeRestCatalogGetCredentials, "Number of 'get credentials' requests to Iceberg REST catalog.", ValueType::Number) \
M(DataLakeRestCatalogGetCredentialsMicroseconds, "Total time of 'get credentials' requests to Iceberg REST catalog.", ValueType::Microseconds) \
M(DataLakeRestCatalogAuthTokenCachedValid, "Number of requests to Iceberg REST catalog that reused a cached access token and did not fetch a new one.", ValueType::Number) \
M(DataLakeRestCatalogAuthTokenRetrieve, "Number of new access tokens fetched for Iceberg REST catalog (OAuth client-credentials or GCP metadata/ADC).", ValueType::Number) \
M(DataLakeRestCatalogAuthTokenRefreshedMicroseconds, "Total time spent fetching access tokens for Iceberg REST catalog.", ValueType::Microseconds) \
M(DataLakeRestCatalogUnauthorized, "Number of Iceberg REST catalog HTTP requests retried with a new access token after HTTP 401 or 403.", ValueType::Number) \
M(DataLakeRestCatalogCreateNamespace, "Number of 'create namespace' requests to Iceberg REST catalog.", ValueType::Number) \
M(DataLakeRestCatalogCreateNamespaceMicroseconds, "Total time of 'create namespace' requests to Iceberg REST catalog.", ValueType::Microseconds) \
M(DataLakeRestCatalogCreateTable, "Number of 'create table' requests to Iceberg REST catalog.", ValueType::Number) \
M(DataLakeRestCatalogCreateTableMicroseconds, "Total time of 'create table' requests to Iceberg REST catalog.", ValueType::Microseconds) \
M(DataLakeRestCatalogUpdateTable, "Number of 'update table' requests to Iceberg REST catalog.", ValueType::Number) \
M(DataLakeRestCatalogUpdateTableMicroseconds, "Total time of 'update table' requests to Iceberg REST catalog.", ValueType::Microseconds) \
M(DataLakeRestCatalogDropTable, "Number of 'drop table' requests to Iceberg REST catalog.", ValueType::Number) \
M(DataLakeRestCatalogDropTableMicroseconds, "Total time of 'drop table' requests to Iceberg REST catalog.", ValueType::Microseconds) \
\
M(DataLakeGlueCatalogGetDatabases, "Number of 'get databases' requests to Iceberg Glue catalog.", ValueType::Number) \
M(DataLakeGlueCatalogGetDatabasesMicroseconds, "Total time of 'get databases' requests to Iceberg Glue catalog.", ValueType::Microseconds) \
M(DataLakeGlueCatalogGetTables, "Number of 'get tables' requests to Iceberg Glue catalog.", ValueType::Number) \
M(DataLakeGlueCatalogGetTablesMicroseconds, "Total time of 'get tables' requests to Iceberg Glue catalog.", ValueType::Microseconds) \
M(DataLakeGlueCatalogGetTable, "Number of 'get table' requests to Iceberg Glue catalog.", ValueType::Number) \
M(DataLakeGlueCatalogGetTableMicroseconds, "Total time of 'get table' requests to Iceberg Glue catalog.", ValueType::Microseconds) \
M(DataLakeGlueCatalogCreateDatabase, "Number of 'create database' requests to Iceberg Glue catalog.", ValueType::Number) \
M(DataLakeGlueCatalogCreateDatabaseMicroseconds, "Total time of 'create database' requests to Iceberg Glue catalog.", ValueType::Microseconds) \
M(DataLakeGlueCatalogCreateTable, "Number of 'create table' requests to Iceberg Glue catalog.", ValueType::Number) \
M(DataLakeGlueCatalogCreateTableMicroseconds, "Total time of 'create table' requests to Iceberg Glue catalog.", ValueType::Microseconds) \
M(DataLakeGlueCatalogUpdateTable, "Number of 'update table' requests to Iceberg Glue catalog.", ValueType::Number) \
M(DataLakeGlueCatalogUpdateTableMicroseconds, "Total time of 'update table' requests to Iceberg Glue catalog.", ValueType::Microseconds) \
M(DataLakeGlueCatalogDropTable, "Number of 'drop table' requests to Iceberg Glue catalog.", ValueType::Number) \
M(DataLakeGlueCatalogDropTableMicroseconds, "Total time of 'drop table' requests to Iceberg Glue catalog.", ValueType::Microseconds) \
\
M(DataLakeUnityCatalogGetTables, "Number of 'get tables' requests to Iceberg Unity catalog.", ValueType::Number) \
M(DataLakeUnityCatalogGetTablesMicroseconds, "Total time of 'get tables' requests to Iceberg Unity catalog.", ValueType::Microseconds) \
M(DataLakeUnityCatalogGetTable, "Number of 'get table' requests to Iceberg Unity catalog.", ValueType::Number) \
M(DataLakeUnityCatalogGetTableMicroseconds, "Total time of 'get table' requests to Iceberg Unity catalog.", ValueType::Microseconds) \
M(DataLakeUnityCatalogGetTableMetadata, "Number of 'get table metadata' requests to Iceberg Unity catalog.", ValueType::Number) \
M(DataLakeUnityCatalogGetTableMetadataMicroseconds, "Total time of 'get table metadata' requests to Iceberg Unity catalog.", ValueType::Microseconds) \
M(DataLakeUnityCatalogGetSchemas, "Number of 'get schemas' requests to Iceberg Unity catalog.", ValueType::Number) \
M(DataLakeUnityCatalogGetSchemasMicroseconds, "Total time of 'get schemas' requests to Iceberg Unity catalog.", ValueType::Microseconds) \
M(DataLakeUnityCatalogGetCredentials, "Number of 'get credentials' requests to Iceberg Unity catalog.", ValueType::Number) \
M(DataLakeUnityCatalogGetCredentialsMicroseconds, "Total time of 'get credentials' requests to Iceberg Unity catalog.", ValueType::Microseconds) \
\


#ifdef APPLY_FOR_EXTERNAL_EVENTS
#define APPLY_FOR_EVENTS(M) APPLY_FOR_BUILTIN_EVENTS(M) APPLY_FOR_EXTERNAL_EVENTS(M)
Expand Down
3 changes: 3 additions & 0 deletions src/Core/Settings.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -9151,6 +9151,9 @@ Multiple algorithms can be specified as a comma-separated list, e.g. `dphyp,gree
)", EXPERIMENTAL) \
DECLARE(Bool, allow_experimental_database_paimon_rest_catalog, false, R"(
Allow experimental database engine DataLakeCatalog with catalog_type = 'paimon_rest'
)", EXPERIMENTAL) \
DECLARE(Bool, allow_experimental_database_s3_tables, false, R"(
Allow experimental database engine DataLakeCatalog with catalog_type = 's3tables' (Amazon S3 Tables Iceberg REST with SigV4)
)", EXPERIMENTAL) \
DECLARE(UInt64, webassembly_udf_max_fuel, 100'000, R"(
Fuel limit per WebAssembly UDF instance execution. Each WebAssembly instruction consumes some amount of fuel. The value is scaled by 1024 before being passed to the runtime, so `webassembly_udf_max_fuel = 1` corresponds to approximately 1024 fuel units. Set to 0 for no finite limit. Applies only to functions whose per-function setting `webassembly_udf_enable_fuel` is true, which is the default.
Expand Down
1 change: 1 addition & 0 deletions src/Core/SettingsChangesHistory.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -265,6 +265,7 @@ const VersionToSettingsChangesMap & getSettingsChangesHistory()
{"allow_experimental_query_deduplication", false, false, "The setting is obsolete, the feature has been removed."},
{"query_plan_min_columns_for_join_lazy_indexing", 0, 3, "Control the minimum number of payload columns from the left side required for enabling lazy indexing optimization in JOIN"},
{"query_plan_max_limit_for_join_lazy_indexing", 1000, 1000, "Added new setting to control maximum limit value that allows to use query plan for lazy join indexing optimization. If zero, there is no limit"},
{"allow_experimental_database_s3_tables", false, false, "New setting to enable experimental database S3 tables (AWS Iceberg REST catalog)."},
{"statistics_max_set_size_for_exact_selectivity_estimation", 10000, 10000, "The bound on the cost of estimating the selectivity of `IN` with a large set is kept under `compatibility` with an earlier version: the previous value is deliberately equal to the new one, so that the uncapped estimation, which could add hundreds of milliseconds to the planning of a single query, is not restored."},
});

Expand Down
7 changes: 4 additions & 3 deletions src/Databases/DataLake/DatabaseDataLake.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -104,6 +104,7 @@ namespace Setting
extern const SettingsBool allow_experimental_database_glue_catalog;
extern const SettingsBool allow_experimental_database_hms_catalog;
extern const SettingsBool allow_experimental_database_paimon_rest_catalog;
extern const SettingsBool allow_experimental_database_s3_tables;
extern const SettingsBool use_hive_partitioning;
extern const SettingsBool log_queries;
extern const SettingsBool parallel_replicas_for_cluster_engines;
Expand Down Expand Up @@ -1664,11 +1665,11 @@ void registerDatabaseDataLake(DatabaseFactory & factory)
case DatabaseDataLakeCatalogType::S3_TABLES:
{
if (!args.create_query.attach
&& !args.context->getSettingsRef()[Setting::allow_experimental_database_iceberg])
&& !args.context->getSettingsRef()[Setting::allow_experimental_database_s3_tables])
{
throw Exception(ErrorCodes::SUPPORT_IS_DISABLED,
"DatabaseDataLake with S3 Tables catalog (Iceberg REST) is beta. "
"To allow its usage, enable setting allow_database_iceberg");
"DatabaseDataLake with S3 Tables catalog is experimental. "
"To allow its usage, enable setting allow_experimental_database_s3_tables");
}

engine_func->name = "Iceberg";
Expand Down
68 changes: 62 additions & 6 deletions src/Databases/DataLake/GlueCatalog.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@

#include <Common/Exception.h>
#include <Common/CurrentMetrics.h>
#include <Common/CurrentThread.h>
#include <Core/Settings.h>
#include <Interpreters/Context.h>

Expand Down Expand Up @@ -83,6 +84,24 @@ namespace DB::ServerSetting
extern const ServerSettingsUInt64 s3_retry_attempts;
}

namespace ProfileEvents
{
extern const Event DataLakeGlueCatalogGetDatabases;
extern const Event DataLakeGlueCatalogGetDatabasesMicroseconds;
extern const Event DataLakeGlueCatalogGetTables;
extern const Event DataLakeGlueCatalogGetTablesMicroseconds;
extern const Event DataLakeGlueCatalogGetTable;
extern const Event DataLakeGlueCatalogGetTableMicroseconds;
extern const Event DataLakeGlueCatalogCreateDatabase;
extern const Event DataLakeGlueCatalogCreateDatabaseMicroseconds;
extern const Event DataLakeGlueCatalogCreateTable;
extern const Event DataLakeGlueCatalogCreateTableMicroseconds;
extern const Event DataLakeGlueCatalogUpdateTable;
extern const Event DataLakeGlueCatalogUpdateTableMicroseconds;
extern const Event DataLakeGlueCatalogDropTable;
extern const Event DataLakeGlueCatalogDropTableMicroseconds;
}

namespace CurrentMetrics
{
extern const Metric MarkCacheBytes;
Expand Down Expand Up @@ -212,7 +231,14 @@ DataLake::ICatalog::Namespaces GlueCatalog::getDatabases(const std::string & pre
do
{
request.SetNextToken(next_token);
auto outcome = glue_client->GetDatabases(request);

Aws::Glue::Model::GetDatabasesOutcome outcome;
{
ProfileEvents::increment(ProfileEvents::DataLakeGlueCatalogGetDatabases);
auto timer = DB::CurrentThread::getProfileEvents().timer(ProfileEvents::DataLakeGlueCatalogGetDatabasesMicroseconds);
outcome = glue_client->GetDatabases(request);
}

if (outcome.IsSuccess())
{
const auto & databases_result = outcome.GetResult();
Expand Down Expand Up @@ -261,7 +287,12 @@ CatalogTables GlueCatalog::getTablesForDatabase(const std::string & db_name, siz
do
{
request.SetNextToken(next_token);
auto outcome = glue_client->GetTables(request);
Aws::Glue::Model::GetTablesOutcome outcome;
{
ProfileEvents::increment(ProfileEvents::DataLakeGlueCatalogGetTables);
auto timer = DB::CurrentThread::getProfileEvents().timer(ProfileEvents::DataLakeGlueCatalogGetTablesMicroseconds);
outcome = glue_client->GetTables(request);
}
if (outcome.IsSuccess())
{
const auto & tables_result = outcome.GetResult();
Expand Down Expand Up @@ -339,7 +370,12 @@ bool GlueCatalog::tryGetTableMetadata(
request.SetDatabaseName(database_name);
request.SetName(table_name);

auto outcome = glue_client->GetTable(request);
Aws::Glue::Model::GetTableOutcome outcome;
{
ProfileEvents::increment(ProfileEvents::DataLakeGlueCatalogGetTable);
auto timer = DB::CurrentThread::getProfileEvents().timer(ProfileEvents::DataLakeGlueCatalogGetTableMicroseconds);
outcome = glue_client->GetTable(request);
}
if (outcome.IsSuccess())
{
const auto & table_outcome = outcome.GetResult().GetTable();
Expand Down Expand Up @@ -635,6 +671,8 @@ void GlueCatalog::createNamespaceIfNotExists(const String & namespace_name, cons
db_input.SetName(namespace_name);
create_request.SetDatabaseInput(db_input);

ProfileEvents::increment(ProfileEvents::DataLakeGlueCatalogCreateDatabase);
auto timer = DB::CurrentThread::getProfileEvents().timer(ProfileEvents::DataLakeGlueCatalogCreateDatabaseMicroseconds);
auto outcome = glue_client->CreateDatabase(create_request);
if (!outcome.IsSuccess() && outcome.GetError().GetErrorType() != Aws::Glue::GlueErrors::ALREADY_EXISTS)
{
Expand Down Expand Up @@ -672,7 +710,13 @@ void GlueCatalog::createTable(const String & namespace_name, const String & tabl

request.SetTableInput(table_input);

auto response = glue_client->CreateTable(request);
Aws::Glue::Model::CreateTableOutcome response;

{
ProfileEvents::increment(ProfileEvents::DataLakeGlueCatalogCreateTable);
auto timer = DB::CurrentThread::getProfileEvents().timer(ProfileEvents::DataLakeGlueCatalogCreateTableMicroseconds);
response = glue_client->CreateTable(request);
}

if (!response.IsSuccess())
throw DB::Exception(DB::ErrorCodes::DATALAKE_DATABASE_ERROR, "Can not create metadata in glue catalog: {}", response.GetError().GetMessage());
Expand Down Expand Up @@ -707,7 +751,13 @@ bool GlueCatalog::updateMetadata(const String & namespace_name, const String & t

request.SetTableInput(table_input);

auto response = glue_client->UpdateTable(request);
Aws::Glue::Model::UpdateTableOutcome response;

{
ProfileEvents::increment(ProfileEvents::DataLakeGlueCatalogUpdateTable);
auto timer = DB::CurrentThread::getProfileEvents().timer(ProfileEvents::DataLakeGlueCatalogUpdateTableMicroseconds);
response = glue_client->UpdateTable(request);
}

if (!response.IsSuccess())
throw DB::Exception(DB::ErrorCodes::DATALAKE_DATABASE_ERROR, "Can not update metadata in glue catalog {}", response.GetError().GetMessage());
Expand All @@ -731,7 +781,13 @@ void GlueCatalog::dropTable(const String & namespace_name, const String & table_
request.SetDatabaseName(namespace_name);
request.SetName(table_name);

auto response = glue_client->DeleteTable(request);
Aws::Glue::Model::DeleteTableOutcome response;

{
ProfileEvents::increment(ProfileEvents::DataLakeGlueCatalogDropTable);
auto timer = DB::CurrentThread::getProfileEvents().timer(ProfileEvents::DataLakeGlueCatalogDropTableMicroseconds);
response = glue_client->DeleteTable(request);
}

if (!response.IsSuccess())
throw DB::Exception(
Expand Down
18 changes: 16 additions & 2 deletions src/Databases/DataLake/ICatalog.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -320,8 +320,22 @@ std::string TableMetadata::getMetadataLocation(const std::string & iceberg_metad
metadata_location = metadata_location.substr(storage_type_str.size());
if (data_location.starts_with(storage_type_str))
data_location = data_location.substr(storage_type_str.size());
else if (!endpoint.empty() && data_location.starts_with(endpoint))
data_location = data_location.substr(endpoint.size());
else if (!endpoint.empty())
{
std::string normalized_endpoint = endpoint;
if (normalized_endpoint.ends_with('/'))
normalized_endpoint.pop_back();

if (data_location.starts_with(normalized_endpoint))
{
data_location = data_location.substr(normalized_endpoint.size());
/// `metadata_location` is relative to the bucket (the `s3://` prefix is stripped above),
/// while `data_location` still has the leading slash left over from the endpoint,
/// e.g. "/bucket/table-uuid/". Drop it so that the prefix comparison below works.
if (azure_account_with_suffix.empty() && data_location.starts_with('/'))
data_location = data_location.substr(1);
}
}

if (metadata_location.starts_with(data_location))
{
Expand Down
Loading
Loading