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
1 change: 1 addition & 0 deletions CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -451,6 +451,7 @@ set(DUCKDB_SRC_FILES
src/duckdb/ub_src_optimizer_join_order.cpp
src/duckdb/ub_src_optimizer_pullup.cpp
src/duckdb/ub_src_optimizer_pushdown.cpp
src/duckdb/ub_src_optimizer_relation_statistics.cpp
src/duckdb/ub_src_optimizer_rule.cpp
src/duckdb/ub_src_optimizer_statistics_expression.cpp
src/duckdb/ub_src_optimizer_statistics_operator.cpp
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -71,7 +71,7 @@ struct HistogramBinState {
}
}

counts->resize(bin_list.length + 1);
counts->resize(bin_boundaries->size() + 1);
}
};

Expand Down
3 changes: 2 additions & 1 deletion src/duckdb/extension/icu/icu-table-range.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -148,7 +148,7 @@ struct ICUTableRange {

template <bool GENERATE_SERIES>
static unique_ptr<FunctionData> Bind(ClientContext &context, TableFunctionBindInput &input,
vector<LogicalType> &return_types, vector<string> &names) {
vector<LogicalType> &return_types, vector<Identifier> &names) {
auto result = make_uniq<ICURangeBindData>(context, input.inputs);

return_types.push_back(LogicalType::TIMESTAMP_TZ);
Expand Down Expand Up @@ -229,6 +229,7 @@ struct ICUTableRange {
nullptr, Bind<false>, nullptr, RangeDateTimeLocalInit);
range_function.in_out_function = ICUTableRangeFunction<false>;
range_function.cardinality = Cardinality;
range_function.return_type = TableFunctionReturnType::SET_RETURNING_FUNCTION;
range.AddFunction(range_function);

loader.RegisterFunction(range);
Expand Down
2 changes: 1 addition & 1 deletion src/duckdb/extension/icu/icu-timezone.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ struct ICUTimeZoneData : public GlobalTableFunctionState {
};

static duckdb::unique_ptr<FunctionData> ICUTimeZoneBind(ClientContext &context, TableFunctionBindInput &input,
vector<LogicalType> &return_types, vector<string> &names) {
vector<LogicalType> &return_types, vector<Identifier> &names) {
names.emplace_back("name");
return_types.emplace_back(LogicalType::VARCHAR);
names.emplace_back("abbrev");
Expand Down
2 changes: 1 addition & 1 deletion src/duckdb/extension/icu/icu_extension.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -365,7 +365,7 @@ struct ICUCalendarData : public GlobalTableFunctionState {
};

static duckdb::unique_ptr<FunctionData> ICUCalendarBind(ClientContext &context, TableFunctionBindInput &input,
vector<LogicalType> &return_types, vector<string> &names) {
vector<LogicalType> &return_types, vector<Identifier> &names) {
names.emplace_back("name");
return_types.emplace_back(LogicalType::VARCHAR);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -282,7 +282,7 @@ struct ExecuteSqlTableFunction {
};

static unique_ptr<FunctionData> Bind(ClientContext &context, TableFunctionBindInput &input,
vector<LogicalType> &return_types, vector<string> &names) {
vector<LogicalType> &return_types, vector<Identifier> &names) {
JSONFunctionLocalState local_state(context);
auto alc = local_state.json_allocator->GetYYAlc();

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ namespace duckdb {
enum class JSONTableInOutType { EACH, TREE };

static unique_ptr<FunctionData> JSONTableInOutBind(ClientContext &, TableFunctionBindInput &input,
vector<LogicalType> &return_types, vector<string> &names) {
vector<LogicalType> &return_types, vector<Identifier> &names) {
const child_list_t<LogicalType> schema {
{"key", LogicalType::VARCHAR}, {"value", LogicalType::JSON()}, {"type", LogicalType::VARCHAR},
{"atom", LogicalType::JSON()}, {"id", LogicalType::UBIGINT}, {"parent", LogicalType::UBIGINT},
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,9 +42,9 @@ class ExpressionColumnReader : public ColumnReader {
static constexpr const PhysicalType TYPE = PhysicalType::INVALID;

public:
ExpressionColumnReader(ClientContext &context, vector<unique_ptr<ColumnReader>> child_readers,
ExpressionColumnReader(ClientContext &context_p, vector<unique_ptr<ColumnReader>> child_readers,
unique_ptr<Expression> expr, const ParquetColumnSchema &schema);
ExpressionColumnReader(ClientContext &context, vector<unique_ptr<ColumnReader>> child_readers,
ExpressionColumnReader(ClientContext &context_p, vector<unique_ptr<ColumnReader>> child_readers,
unique_ptr<Expression> expr, unique_ptr<ParquetColumnSchema> owned_schema);

//! Reader(s) to produce the input(s) for the expression
Expand All @@ -61,6 +61,7 @@ class ExpressionColumnReader : public ColumnReader {
TProtocol &protocol_p) override;

idx_t Read(ColumnReaderInput &input, Vector &result) override;
unique_ptr<BaseStatistics> Stats(idx_t row_group_idx_p, const vector<ColumnChunk> &columns) override;

void Select(ColumnReaderInput &input, Vector &result, const SelectionVector &sel,
idx_t approved_tuple_count) override;
Expand All @@ -87,6 +88,7 @@ class ExpressionColumnReader : public ColumnReader {
}

private:
ClientContext &context;
void InitializeChunk();
};

Expand Down
34 changes: 17 additions & 17 deletions src/duckdb/extension/parquet/parquet_metadata.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -118,7 +118,7 @@ class ParquetMetaDataOperator {
public:
template <ParquetMetadataOperatorType OP_TYPE>
static unique_ptr<FunctionData> Bind(ClientContext &context, TableFunctionBindInput &input,
vector<LogicalType> &return_types, vector<string> &names);
vector<LogicalType> &return_types, vector<Identifier> &names);
static unique_ptr<GlobalTableFunctionState> InitGlobal(ClientContext &context, TableFunctionInitInput &input);
template <ParquetMetadataOperatorType OP_TYPE>
static unique_ptr<LocalTableFunctionState> InitLocal(ExecutionContext &context, TableFunctionInitInput &input,
Expand All @@ -128,7 +128,7 @@ class ParquetMetaDataOperator {
const GlobalTableFunctionState *global_state);

template <ParquetMetadataOperatorType OP_TYPE>
static void BindSchema(vector<LogicalType> &return_types, vector<string> &names);
static void BindSchema(vector<LogicalType> &return_types, vector<Identifier> &names);

static OperatorPartitionData GetPartitionData(ClientContext &context, TableFunctionGetPartitionInput &input);
};
Expand Down Expand Up @@ -252,7 +252,7 @@ class ParquetRowGroupMetadataProcessor : public ParquetMetadataFileProcessor {

template <>
void ParquetMetaDataOperator::BindSchema<ParquetMetadataOperatorType::META_DATA>(vector<LogicalType> &return_types,
vector<string> &names) {
vector<Identifier> &names) {
names.emplace_back("file_name");
return_types.emplace_back(LogicalType::VARCHAR);

Expand Down Expand Up @@ -539,7 +539,7 @@ class ParquetSchemaProcessor : public ParquetMetadataFileProcessor {

template <>
void ParquetMetaDataOperator::BindSchema<ParquetMetadataOperatorType::SCHEMA>(vector<LogicalType> &return_types,
vector<string> &names) {
vector<Identifier> &names) {
names.emplace_back("file_name");
return_types.emplace_back(LogicalType::VARCHAR);

Expand Down Expand Up @@ -688,7 +688,7 @@ class ParquetKeyValueMetadataProcessor : public ParquetMetadataFileProcessor {

template <>
void ParquetMetaDataOperator::BindSchema<ParquetMetadataOperatorType::KEY_VALUE_META_DATA>(
vector<LogicalType> &return_types, vector<string> &names) {
vector<LogicalType> &return_types, vector<Identifier> &names) {
names.emplace_back("file_name");
return_types.emplace_back(LogicalType::VARCHAR);

Expand Down Expand Up @@ -725,7 +725,7 @@ class ParquetFileMetadataProcessor : public ParquetMetadataFileProcessor {

template <>
void ParquetMetaDataOperator::BindSchema<ParquetMetadataOperatorType::FILE_META_DATA>(vector<LogicalType> &return_types,
vector<string> &names) {
vector<Identifier> &names) {
names.emplace_back("file_name");
return_types.emplace_back(LogicalType::VARCHAR);

Expand Down Expand Up @@ -822,7 +822,7 @@ class ParquetBloomProbeProcessor : public ParquetMetadataFileProcessor {

template <>
void ParquetMetaDataOperator::BindSchema<ParquetMetadataOperatorType::BLOOM_PROBE>(vector<LogicalType> &return_types,
vector<string> &names) {
vector<Identifier> &names) {
names.emplace_back("file_name");
return_types.emplace_back(LogicalType::VARCHAR);

Expand Down Expand Up @@ -942,44 +942,44 @@ void FullMetadataProcessor::PopulateMetadata(ParquetMetadataFileProcessor &proce

template <>
void ParquetMetaDataOperator::BindSchema<ParquetMetadataOperatorType::FULL_METADATA>(vector<LogicalType> &return_types,
vector<string> &names) {
vector<Identifier> &names) {
names.emplace_back("parquet_file_metadata");
vector<LogicalType> file_meta_types;
vector<string> file_meta_names;
vector<Identifier> file_meta_names;
ParquetMetaDataOperator::BindSchema<ParquetMetadataOperatorType::FILE_META_DATA>(file_meta_types, file_meta_names);
child_list_t<LogicalType> file_meta_children;
for (idx_t i = 0; i < file_meta_types.size(); i++) {
file_meta_children.emplace_back(make_pair(file_meta_names[i], file_meta_types[i]));
file_meta_children.emplace_back(make_pair(file_meta_names[i].GetIdentifierName(), file_meta_types[i]));
}
return_types.emplace_back(LogicalType::LIST(LogicalType::STRUCT(std::move(file_meta_children))));

names.emplace_back("parquet_metadata");
vector<LogicalType> row_group_types;
vector<string> row_group_names;
vector<Identifier> row_group_names;
ParquetMetaDataOperator::BindSchema<ParquetMetadataOperatorType::META_DATA>(row_group_types, row_group_names);
child_list_t<LogicalType> row_group_children;
for (idx_t i = 0; i < row_group_types.size(); i++) {
row_group_children.emplace_back(make_pair(row_group_names[i], row_group_types[i]));
row_group_children.emplace_back(make_pair(row_group_names[i].GetIdentifierName(), row_group_types[i]));
}
return_types.emplace_back(LogicalType::LIST(LogicalType::STRUCT(std::move(row_group_children))));

names.emplace_back("parquet_schema");
vector<LogicalType> schema_types;
vector<string> schema_names;
vector<Identifier> schema_names;
ParquetMetaDataOperator::BindSchema<ParquetMetadataOperatorType::SCHEMA>(schema_types, schema_names);
child_list_t<LogicalType> schema_children;
for (idx_t i = 0; i < schema_types.size(); i++) {
schema_children.emplace_back(make_pair(schema_names[i], schema_types[i]));
schema_children.emplace_back(make_pair(schema_names[i].GetIdentifierName(), schema_types[i]));
}
return_types.emplace_back(LogicalType::LIST(LogicalType::STRUCT(std::move(schema_children))));

names.emplace_back("parquet_kv_metadata");
vector<LogicalType> kv_types;
vector<string> kv_names;
vector<Identifier> kv_names;
ParquetMetaDataOperator::BindSchema<ParquetMetadataOperatorType::KEY_VALUE_META_DATA>(kv_types, kv_names);
child_list_t<LogicalType> kv_children;
for (idx_t i = 0; i < kv_types.size(); i++) {
kv_children.emplace_back(make_pair(kv_names[i], kv_types[i]));
kv_children.emplace_back(make_pair(kv_names[i].GetIdentifierName(), kv_types[i]));
}
return_types.emplace_back(LogicalType::LIST(LogicalType::STRUCT(std::move(kv_children))));
}
Expand Down Expand Up @@ -1008,7 +1008,7 @@ void FullMetadataProcessor::ReadRow(vector<reference<Vector>> &output, idx_t row

template <ParquetMetadataOperatorType OP_TYPE>
unique_ptr<FunctionData> ParquetMetaDataOperator::Bind(ClientContext &context, TableFunctionBindInput &input,
vector<LogicalType> &return_types, vector<string> &names) {
vector<LogicalType> &return_types, vector<Identifier> &names) {
// Extract file paths from input using MultiFileReader (handles both single files and arrays)
auto multi_file_reader = MultiFileReader::CreateDefault("ParquetMetadata");
auto glob_input = FileGlobInput(FileGlobOptions::FALLBACK_GLOB, "parquet");
Expand Down
2 changes: 1 addition & 1 deletion src/duckdb/extension/parquet/parquet_multi_file_info.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -304,7 +304,7 @@ static unique_ptr<FunctionData> ParquetScanDeserialize(Deserializer &deserialize
auto &context = deserializer.Get<ClientContext &>();
auto files = deserializer.ReadProperty<vector<string>>(100, "files");
auto types = deserializer.ReadProperty<vector<LogicalType>>(101, "types");
auto names = deserializer.ReadProperty<vector<string>>(102, "names");
auto names = StringsToIdentifiers(deserializer.ReadProperty<vector<string>>(102, "names"));
auto serialization = deserializer.ReadProperty<ParquetOptionsSerialization>(103, "parquet_options");
auto table_columns =
deserializer.ReadPropertyWithExplicitDefault<vector<string>>(104, "table_columns", vector<string> {});
Expand Down
11 changes: 4 additions & 7 deletions src/duckdb/extension/parquet/parquet_reader.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -1677,13 +1677,10 @@ void ParquetReader::PrepareRowGroupBuffer(ClientContext &context, ParquetReaderS
has_min_max = group.columns[schema_column_index].meta_data.statistics.__isset.min_value &&
group.columns[schema_column_index].meta_data.statistics.__isset.max_value;
}
if (is_expression) {
// no pruning possible for expressions
prune_result = FilterPropagateResult::NO_PRUNING_POSSIBLE;
} else if (!is_generated_column && has_min_max &&
(column_reader.Type().id() == LogicalTypeId::FLOAT ||
column_reader.Type().id() == LogicalTypeId::DOUBLE) &&
parquet_options.can_have_nan) {
if (!is_expression && !is_generated_column && has_min_max &&
(column_reader.Type().id() == LogicalTypeId::FLOAT ||
column_reader.Type().id() == LogicalTypeId::DOUBLE) &&
parquet_options.can_have_nan) {
// floating point columns can have NaN values in addition to the min/max bounds defined in the file
// in order to do optimal pruning - we prune based on the [min, max] of the file followed by pruning
// based on nan
Expand Down
29 changes: 25 additions & 4 deletions src/duckdb/extension/parquet/reader/expression_column_reader.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
#include "parquet_reader.hpp"
#include "duckdb/common/types/vector.hpp"
#include "duckdb/common/vector/flat_vector.hpp"
#include "duckdb/planner/filter/expression_filter.hpp"

namespace duckdb_apache {
namespace thrift {
Expand All @@ -23,10 +24,11 @@ class ClientContext;
//===--------------------------------------------------------------------===//
// Expression Column Reader
//===--------------------------------------------------------------------===//
ExpressionColumnReader::ExpressionColumnReader(ClientContext &context, vector<unique_ptr<ColumnReader>> child_readers_p,
ExpressionColumnReader::ExpressionColumnReader(ClientContext &context_p,
vector<unique_ptr<ColumnReader>> child_readers_p,
unique_ptr<Expression> expr_p, const ParquetColumnSchema &schema_p)
: ColumnReader(child_readers_p[0]->Reader(), schema_p), child_readers(std::move(child_readers_p)),
expr(std::move(expr_p)), executor(context, expr.get()) {
expr(std::move(expr_p)), executor(context_p, expr.get()), context(context_p) {
if (child_readers.empty()) {
throw InternalException("Can't instantiate an ExpressionColumnReader with 0 children");
}
Expand All @@ -38,11 +40,13 @@ ExpressionColumnReader::ExpressionColumnReader(ClientContext &context, vector<un
InitializeChunk();
}

ExpressionColumnReader::ExpressionColumnReader(ClientContext &context, vector<unique_ptr<ColumnReader>> child_readers_p,
ExpressionColumnReader::ExpressionColumnReader(ClientContext &context_p,
vector<unique_ptr<ColumnReader>> child_readers_p,
unique_ptr<Expression> expr_p,
unique_ptr<ParquetColumnSchema> owned_schema_p)
: ColumnReader(child_readers_p[0]->Reader(), *owned_schema_p), child_readers(std::move(child_readers_p)),
expr(std::move(expr_p)), executor(context, expr.get()), owned_schema(std::move(owned_schema_p)) {
expr(std::move(expr_p)), executor(context_p, expr.get()), owned_schema(std::move(owned_schema_p)),
context(context_p) {
if (child_readers.empty()) {
throw InternalException("Can't instantiate an ExpressionColumnReader with 0 children");
}
Expand All @@ -69,6 +73,23 @@ void ExpressionColumnReader::InitializeRead(idx_t row_group_idx_p, idx_t row_gro
}
}

unique_ptr<BaseStatistics> ExpressionColumnReader::Stats(idx_t row_group_idx_p, const vector<ColumnChunk> &columns) {
if (Schema().schema_type != ParquetColumnSchemaType::EXPRESSION) {
return ColumnReader::Stats(row_group_idx_p, columns);
}
vector<BaseStatistics> input_stats;
input_stats.reserve(child_readers.size());
for (auto &child_reader : child_readers) {
auto child_stats = child_reader->Stats(row_group_idx_p, columns);
if (!child_stats) {
return nullptr;
}
input_stats.push_back(child_stats->Copy());
}
return ExpressionFilter::TryGetExpressionStatistics(
context, *expr, array_ptr<const BaseStatistics>(input_stats.data(), input_stats.size()));
}

static void ReverseSelectionVector(const SelectionVector &input, SelectionVector &output, idx_t input_count,
idx_t result_count) {
//! For an input selection vector: [5, 10],
Expand Down
Loading
Loading