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 @@ -481,6 +481,7 @@ set(DUCKDB_SRC_FILES
src/duckdb/ub_src_planner_expression_binder.cpp
src/duckdb/ub_src_planner_filter.cpp
src/duckdb/ub_src_planner_operator.cpp
src/duckdb/ub_src_planner_sql_export.cpp
src/duckdb/ub_src_planner_subquery.cpp
src/duckdb/ub_src_storage.cpp
src/duckdb/ub_src_storage_buffer.cpp
Expand Down
15 changes: 15 additions & 0 deletions src/duckdb/extension/core_functions/scalar/date/date_part.cpp
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
#include "duckdb/parser/expression/function_expression.hpp"
#include "duckdb/common/vector/struct_vector.hpp"
#include "core_functions/scalar/date_functions.hpp"
#include "duckdb/common/case_insensitive_map.hpp"
Expand Down Expand Up @@ -2716,9 +2717,23 @@ ScalarFunctionSet JulianDayFun::GetFunctions() {
return operator_set;
}

//! Binding date_part with a constant part replaces it with the unary part function (year, month, ...), so the bound
//! call has one argument and must be rendered under the replacement's own name
static unique_ptr<ParsedExpression> DatePartUnbind(FunctionUnbindInput &input) {
auto &function = input.expression.Function();
if (input.children.size() == 1) {
return make_uniq<FunctionExpression>(function.GetQualifiedName(), std::move(input.children));
}
if (input.children.size() != 2) {
return nullptr;
}
return make_uniq<FunctionExpression>(function.GetDefinition()->GetQualifiedName(), std::move(input.children));
}

// Names the "part,ts" pair shared by date_part's per-type overloads.
static ScalarFunction NamePartTsArguments(ScalarFunction fun, const LogicalType &type) {
fun.GetSignature().AddParameter("part", LogicalType::VARCHAR).AddParameter("ts", type);
fun.SetUnbindCallback(DatePartUnbind);
return fun;
}

Expand Down
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
#include "duckdb/parser/expression/function_expression.hpp"
#include "duckdb/common/vector/map_vector.hpp"
#include "duckdb/common/vector/struct_vector.hpp"
#include "core_functions/scalar/struct_functions.hpp"
Expand Down Expand Up @@ -140,11 +141,25 @@ static unique_ptr<BaseStatistics> StructUpdateStats(ClientContext &context, Func
return new_stats.ToUnique();
}

static unique_ptr<ParsedExpression> StructUpdateUnbind(FunctionUnbindInput &input) {
vector<FunctionArgument> arguments;
for (idx_t i = 0; i < input.children.size(); i++) {
auto name = i == 0 ? Identifier() : input.expression.GetChildren()[i]->GetAlias();
if (i > 0 && name.empty()) {
return nullptr;
}
arguments.emplace_back(std::move(name), std::move(input.children[i]));
}
return make_uniq<FunctionExpression>(input.expression.Function().GetDefinition()->GetQualifiedName(),
std::move(arguments));
}

ScalarFunction StructUpdateFun::GetFunction() {
ScalarFunction fun({}, LogicalTypeId::STRUCT, StructUpdateFunction, StructUpdateBind, StructUpdateStats);
fun.SetNullHandling(FunctionNullHandling::SPECIAL_HANDLING);
fun.GetSignature().AddParameter("struct", LogicalType::ANY).AddKwargsParameter("kwargs", LogicalType::ANY);
fun.GetProperties().SetRequiresExpressionNames(true);
fun.SetUnbindCallback(StructUpdateUnbind);
fun.SetSerializeCallback(VariableReturnBindData::Serialize);
fun.SetDeserializeCallback(VariableReturnBindData::Deserialize);
return fun;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -333,8 +333,9 @@ TableFunction JSONFunctions::GetJSONTableFunction(Identifier name, shared_ptr<JS
// the schema is determined by combining the schemas of up to 32 files - the keys of the files are unified, so a
// file does not need to have every column of the combined schema
settings.maximum_sample_files = 32;
return TableFunctionMultiFileWrapper::CreateFunction(std::move(single_file_function), std::move(name),
std::move(settings));
auto function = TableFunctionMultiFileWrapper::CreateFunction(std::move(single_file_function), std::move(name),
std::move(settings));
return function;
}

static TableFunctionSet CreateJSONFunctionSet(Identifier name, shared_ptr<JSONScanInfo> function_info) {
Expand Down
30 changes: 14 additions & 16 deletions src/duckdb/extension/parquet/column_reader.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -1012,6 +1012,20 @@ static unique_ptr<ColumnReader> CreateDecimalReader(const ParquetReader &reader,
}
}

static_assert(ParquetTimestampLogicalType(ParquetExtraTypeInfo::UNIT_NS) != LogicalTypeId::TIMESTAMP);
static_assert(ParquetTimestampTzLogicalType(ParquetExtraTypeInfo::UNIT_NS) != LogicalTypeId::TIMESTAMP_TZ);

static_assert(ParquetTimestampLogicalType(ParquetExtraTypeInfo::IMPALA_TIMESTAMP) != LogicalTypeId::TIMESTAMP_NS);
static_assert(ParquetTimestampTzLogicalType(ParquetExtraTypeInfo::IMPALA_TIMESTAMP) != LogicalTypeId::TIMESTAMP_TZ_NS);
static_assert(ParquetTimestampLogicalType(ParquetExtraTypeInfo::UNIT_MS) != LogicalTypeId::TIMESTAMP_NS);
static_assert(ParquetTimestampTzLogicalType(ParquetExtraTypeInfo::UNIT_MS) != LogicalTypeId::TIMESTAMP_TZ_NS);
static_assert(ParquetTimestampLogicalType(ParquetExtraTypeInfo::UNIT_MICROS) != LogicalTypeId::TIMESTAMP_NS);
static_assert(ParquetTimestampTzLogicalType(ParquetExtraTypeInfo::UNIT_MICROS) != LogicalTypeId::TIMESTAMP_TZ_NS);

static_assert(ParquetTimeLogicalType(ParquetExtraTypeInfo::UNIT_NS) != LogicalTypeId::TIME);
static_assert(ParquetTimeLogicalType(ParquetExtraTypeInfo::UNIT_MS) != LogicalTypeId::TIME_NS);
static_assert(ParquetTimeLogicalType(ParquetExtraTypeInfo::UNIT_MICROS) != LogicalTypeId::TIME_NS);

unique_ptr<ColumnReader> ColumnReader::CreateReader(const ParquetReader &reader, const ParquetColumnSchema &schema) {
switch (schema.type.id()) {
case LogicalTypeId::BOOLEAN:
Expand Down Expand Up @@ -1052,22 +1066,12 @@ unique_ptr<ColumnReader> ColumnReader::CreateReader(const ParquetReader &reader,
case ParquetExtraTypeInfo::UNIT_MICROS:
return make_uniq<CallbackColumnReader<int64_t, timestamp_t, ParquetTimestampMicrosToTimestamp>>(reader,
schema);
case ParquetExtraTypeInfo::UNIT_NS:
return make_uniq<CallbackColumnReader<int64_t, timestamp_t, ParquetTimestampNsToTimestamp>>(reader, schema);
default:
throw InternalException("TIMESTAMP requires type info");
}
case LogicalTypeId::TIMESTAMP_NS:
case LogicalTypeId::TIMESTAMP_TZ_NS:
switch (schema.type_info) {
case ParquetExtraTypeInfo::IMPALA_TIMESTAMP:
return make_uniq<CallbackColumnReader<Int96, timestamp_ns_t, ImpalaTimestampToTimestampNS>>(reader, schema);
case ParquetExtraTypeInfo::UNIT_MS:
return make_uniq<CallbackColumnReader<int64_t, timestamp_ns_t, ParquetTimestampMsToTimestampNs>>(reader,
schema);
case ParquetExtraTypeInfo::UNIT_MICROS:
return make_uniq<CallbackColumnReader<int64_t, timestamp_ns_t, ParquetTimestampUsToTimestampNs>>(reader,
schema);
case ParquetExtraTypeInfo::UNIT_NS:
return make_uniq<CallbackColumnReader<int64_t, timestamp_ns_t, ParquetTimestampNsToTimestampNs>>(reader,
schema);
Expand All @@ -1082,17 +1086,11 @@ unique_ptr<ColumnReader> ColumnReader::CreateReader(const ParquetReader &reader,
return make_uniq<CallbackColumnReader<int32_t, dtime_t, ParquetMsIntToTime>>(reader, schema);
case ParquetExtraTypeInfo::UNIT_MICROS:
return make_uniq<CallbackColumnReader<int64_t, dtime_t, ParquetIntToTime>>(reader, schema);
case ParquetExtraTypeInfo::UNIT_NS:
return make_uniq<CallbackColumnReader<int64_t, dtime_t, ParquetNsIntToTime>>(reader, schema);
default:
throw InternalException("TIME requires type info");
}
case LogicalTypeId::TIME_NS:
switch (schema.type_info) {
case ParquetExtraTypeInfo::UNIT_MS:
return make_uniq<CallbackColumnReader<int32_t, dtime_ns_t, ParquetMsIntToTimeNs>>(reader, schema);
case ParquetExtraTypeInfo::UNIT_MICROS:
return make_uniq<CallbackColumnReader<int64_t, dtime_ns_t, ParquetUsIntToTimeNs>>(reader, schema);
case ParquetExtraTypeInfo::UNIT_NS:
return make_uniq<CallbackColumnReader<int64_t, dtime_ns_t, ParquetIntToTimeNs>>(reader, schema);
default:
Expand Down
48 changes: 48 additions & 0 deletions src/duckdb/extension/parquet/include/parquet_column_schema.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,54 @@ enum class ParquetExtraTypeInfo {
FLOAT16
};

constexpr LogicalTypeId ParquetTimestampLogicalType(ParquetExtraTypeInfo type_info) {
switch (type_info) {
case ParquetExtraTypeInfo::IMPALA_TIMESTAMP:
case ParquetExtraTypeInfo::UNIT_MS:
case ParquetExtraTypeInfo::UNIT_MICROS:
return LogicalTypeId::TIMESTAMP;
case ParquetExtraTypeInfo::UNIT_NS:
return LogicalTypeId::TIMESTAMP_NS;
default:
return LogicalTypeId::INVALID;
}
}

constexpr LogicalTypeId ParquetTimestampTzLogicalType(ParquetExtraTypeInfo type_info) {
switch (type_info) {
case ParquetExtraTypeInfo::UNIT_NS:
return LogicalTypeId::TIMESTAMP_TZ_NS;
case ParquetExtraTypeInfo::UNIT_MS:
case ParquetExtraTypeInfo::UNIT_MICROS:
return LogicalTypeId::TIMESTAMP_TZ;
default:
return LogicalTypeId::INVALID;
}
}

constexpr LogicalTypeId ParquetTimeLogicalType(ParquetExtraTypeInfo type_info) {
switch (type_info) {
case ParquetExtraTypeInfo::UNIT_NS:
return LogicalTypeId::TIME_NS;
case ParquetExtraTypeInfo::UNIT_MS:
case ParquetExtraTypeInfo::UNIT_MICROS:
return LogicalTypeId::TIME;
default:
return LogicalTypeId::INVALID;
}
}

constexpr LogicalTypeId ParquetTimeTzLogicalType(ParquetExtraTypeInfo type_info) {
switch (type_info) {
case ParquetExtraTypeInfo::UNIT_MS:
case ParquetExtraTypeInfo::UNIT_MICROS:
case ParquetExtraTypeInfo::UNIT_NS:
return LogicalTypeId::TIME_TZ;
default:
return LogicalTypeId::INVALID;
}
}

struct ParquetColumnSchema {
public:
ParquetColumnSchema() = default;
Expand Down
7 changes: 0 additions & 7 deletions src/duckdb/extension/parquet/include/parquet_timestamp.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -22,24 +22,17 @@ struct Int96 {
};

timestamp_t ImpalaTimestampToTimestamp(const Int96 &raw_ts);
timestamp_ns_t ImpalaTimestampToTimestampNS(const Int96 &raw_ts);
Int96 TimestampToImpalaTimestamp(timestamp_t &ts);

timestamp_t ParquetTimestampMicrosToTimestamp(const int64_t &raw_ts);
timestamp_t ParquetTimestampMsToTimestamp(const int64_t &raw_ts);
timestamp_t ParquetTimestampNsToTimestamp(const int64_t &raw_ts);

timestamp_ns_t ParquetTimestampMsToTimestampNs(const int64_t &raw_ms);
timestamp_ns_t ParquetTimestampUsToTimestampNs(const int64_t &raw_us);
timestamp_ns_t ParquetTimestampNsToTimestampNs(const int64_t &raw_ns);

date_t ParquetIntToDate(const int32_t &raw_date);
dtime_t ParquetMsIntToTime(const int32_t &raw_millis);
dtime_t ParquetIntToTime(const int64_t &raw_micros);
dtime_t ParquetNsIntToTime(const int64_t &raw_nanos);

dtime_ns_t ParquetMsIntToTimeNs(const int32_t &raw_millis);
dtime_ns_t ParquetUsIntToTimeNs(const int64_t &raw_micros);
dtime_ns_t ParquetIntToTimeNs(const int64_t &raw_nanos);

dtime_tz_t ParquetIntToTimeMsTZ(const int32_t &raw_millis);
Expand Down
33 changes: 17 additions & 16 deletions src/duckdb/extension/parquet/parquet_reader.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,14 @@

namespace duckdb {

static_assert(ParquetTimestampTzLogicalType(ParquetExtraTypeInfo::UNIT_MS) == LogicalTypeId::TIMESTAMP_TZ);
static_assert(ParquetTimestampTzLogicalType(ParquetExtraTypeInfo::UNIT_MICROS) == LogicalTypeId::TIMESTAMP_TZ);
static_assert(ParquetTimestampTzLogicalType(ParquetExtraTypeInfo::UNIT_NS) == LogicalTypeId::TIMESTAMP_TZ_NS);

static_assert(ParquetTimeTzLogicalType(ParquetExtraTypeInfo::UNIT_MS) == LogicalTypeId::TIME_TZ);
static_assert(ParquetTimeTzLogicalType(ParquetExtraTypeInfo::UNIT_MICROS) == LogicalTypeId::TIME_TZ);
static_assert(ParquetTimeTzLogicalType(ParquetExtraTypeInfo::UNIT_NS) == LogicalTypeId::TIME_TZ);

const char *ParquetPrefetchStrategyToString(ParquetPrefetchStrategy strategy) {
switch (strategy) {
case ParquetPrefetchStrategy::WHOLE_GROUP:
Expand Down Expand Up @@ -422,14 +430,9 @@ LogicalType ParquetReader::DeriveLogicalType(const SchemaElement &s_ele, const P
throw NotImplementedException("Unimplemented TIMESTAMP encoding - missing UNIT");
}
if (s_ele.logicalType.TIMESTAMP.isAdjustedToUTC) {
if (s_ele.logicalType.TIMESTAMP.unit.__isset.NANOS) {
return LogicalType::TIMESTAMP_TZ_NS;
}
return LogicalType::TIMESTAMP_TZ;
} else if (s_ele.logicalType.TIMESTAMP.unit.__isset.NANOS) {
return LogicalType::TIMESTAMP_NS;
return LogicalType(ParquetTimestampTzLogicalType(schema.type_info));
}
return LogicalType::TIMESTAMP;
return LogicalType(ParquetTimestampLogicalType(schema.type_info));
} else if (s_ele.logicalType.__isset.TIME) {
if (s_ele.logicalType.TIME.unit.__isset.MILLIS) {
schema.type_info = ParquetExtraTypeInfo::UNIT_MS;
Expand All @@ -441,11 +444,9 @@ LogicalType ParquetReader::DeriveLogicalType(const SchemaElement &s_ele, const P
throw NotImplementedException("Unimplemented TIME encoding - missing UNIT");
}
if (s_ele.logicalType.TIME.isAdjustedToUTC) {
return LogicalType::TIME_TZ;
} else if (s_ele.logicalType.TIME.unit.__isset.NANOS) {
return LogicalType::TIME_NS;
return LogicalType(ParquetTimeTzLogicalType(schema.type_info));
}
return LogicalType::TIME;
return LogicalType(ParquetTimeLogicalType(schema.type_info));
}
}
if (s_ele.__isset.converted_type) {
Expand Down Expand Up @@ -511,14 +512,14 @@ LogicalType ParquetReader::DeriveLogicalType(const SchemaElement &s_ele, const P
case ConvertedType::TIMESTAMP_MICROS:
schema.type_info = ParquetExtraTypeInfo::UNIT_MICROS;
if (s_ele.type == Type::INT64) {
return LogicalType::TIMESTAMP;
return LogicalType(ParquetTimestampLogicalType(schema.type_info));
} else {
throw IOException("TIMESTAMP converted type can only be set for value of Type::INT64");
}
case ConvertedType::TIMESTAMP_MILLIS:
schema.type_info = ParquetExtraTypeInfo::UNIT_MS;
if (s_ele.type == Type::INT64) {
return LogicalType::TIMESTAMP;
return LogicalType(ParquetTimestampLogicalType(schema.type_info));
} else {
throw IOException("TIMESTAMP converted type can only be set for value of Type::INT64");
}
Expand Down Expand Up @@ -559,14 +560,14 @@ LogicalType ParquetReader::DeriveLogicalType(const SchemaElement &s_ele, const P
case ConvertedType::TIME_MILLIS:
schema.type_info = ParquetExtraTypeInfo::UNIT_MS;
if (s_ele.type == Type::INT32) {
return LogicalType::TIME;
return LogicalType(ParquetTimeLogicalType(schema.type_info));
} else {
throw IOException("TIME_MILLIS converted type can only be set for value of Type::INT32");
}
case ConvertedType::TIME_MICROS:
schema.type_info = ParquetExtraTypeInfo::UNIT_MICROS;
if (s_ele.type == Type::INT64) {
return LogicalType::TIME;
return LogicalType(ParquetTimeLogicalType(schema.type_info));
} else {
throw IOException("TIME_MICROS converted type can only be set for value of Type::INT64");
}
Expand All @@ -593,7 +594,7 @@ LogicalType ParquetReader::DeriveLogicalType(const SchemaElement &s_ele, const P
return LogicalType::BIGINT;
case Type::INT96: // always a timestamp it would seem
schema.type_info = ParquetExtraTypeInfo::IMPALA_TIMESTAMP;
return LogicalType::TIMESTAMP;
return LogicalType(ParquetTimestampLogicalType(schema.type_info));
case Type::FLOAT:
return LogicalType::FLOAT;
case Type::DOUBLE:
Expand Down
52 changes: 10 additions & 42 deletions src/duckdb/extension/parquet/parquet_statistics.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -253,31 +253,19 @@ Value ParquetStatisticsUtils::ConvertValueInternal(const LogicalType &type, cons
switch (schema_ele.type_info) {
case ParquetExtraTypeInfo::UNIT_MS:
return Value::TIME(Time::FromTimeMs(val));
case ParquetExtraTypeInfo::UNIT_NS:
return Value::TIME(Time::FromTimeNs(val));
case ParquetExtraTypeInfo::UNIT_MICROS:
default:
return Value::TIME(dtime_t(val));
}
}
case LogicalTypeId::TIME_NS: {
int64_t val;
if (stats.size() == sizeof(int32_t)) {
val = Load<int32_t>(stats_data);
} else if (stats.size() == sizeof(int64_t)) {
val = Load<int64_t>(stats_data);
} else {
if (stats.size() != sizeof(int64_t)) {
throw InvalidInputException("Incorrect stats size for type TIME_NS");
}
switch (schema_ele.type_info) {
case ParquetExtraTypeInfo::UNIT_MS:
return Value::TIME_NS(ParquetMsIntToTimeNs(NumericCast<int32_t>(val)));
case ParquetExtraTypeInfo::UNIT_NS:
return Value::TIME_NS(ParquetIntToTimeNs(val));
case ParquetExtraTypeInfo::UNIT_MICROS:
default:
return Value::TIME_NS(dtime_ns_t(val));
if (schema_ele.type_info != ParquetExtraTypeInfo::UNIT_NS) {
throw InternalException("TIME_NS requires nanosecond type info");
}
return Value::TIME_NS(ParquetIntToTimeNs(Load<int64_t>(stats_data)));
}
case LogicalTypeId::TIME_TZ: {
int64_t val;
Expand Down Expand Up @@ -315,9 +303,6 @@ Value ParquetStatisticsUtils::ConvertValueInternal(const LogicalType &type, cons
case ParquetExtraTypeInfo::UNIT_MS:
timestamp_value = ParquetTimestampMsToTimestamp(val);
break;
case ParquetExtraTypeInfo::UNIT_NS:
timestamp_value = ParquetTimestampNsToTimestamp(val);
break;
case ParquetExtraTypeInfo::UNIT_MICROS:
default:
timestamp_value = timestamp_t(val);
Expand All @@ -331,30 +316,13 @@ Value ParquetStatisticsUtils::ConvertValueInternal(const LogicalType &type, cons
}
case LogicalTypeId::TIMESTAMP_TZ_NS:
case LogicalTypeId::TIMESTAMP_NS: {
timestamp_ns_t timestamp_value;
if (schema_ele.type_info == ParquetExtraTypeInfo::IMPALA_TIMESTAMP) {
if (stats.size() != sizeof(Int96)) {
throw InvalidInputException("Incorrect stats size for type TIMESTAMP_NS");
}
timestamp_value = ImpalaTimestampToTimestampNS(Load<Int96>(stats_data));
} else {
if (stats.size() != sizeof(int64_t)) {
throw InvalidInputException("Incorrect stats size for type TIMESTAMP_NS");
}
auto val = Load<int64_t>(stats_data);
switch (schema_ele.type_info) {
case ParquetExtraTypeInfo::UNIT_MS:
timestamp_value = ParquetTimestampMsToTimestampNs(val);
break;
case ParquetExtraTypeInfo::UNIT_NS:
timestamp_value = ParquetTimestampNsToTimestampNs(val);
break;
case ParquetExtraTypeInfo::UNIT_MICROS:
default:
timestamp_value = ParquetTimestampUsToTimestampNs(val);
break;
}
if (stats.size() != sizeof(int64_t)) {
throw InvalidInputException("Incorrect stats size for type TIMESTAMP_NS");
}
if (schema_ele.type_info != ParquetExtraTypeInfo::UNIT_NS) {
throw InternalException("TIMESTAMP_NS requires nanosecond type info");
}
auto timestamp_value = ParquetTimestampNsToTimestampNs(Load<int64_t>(stats_data));
if (type.id() == LogicalTypeId::TIMESTAMP_TZ_NS) {
return Value::TIMESTAMPTZNS(timestamp_tz_ns_t(timestamp_value));
}
Expand Down
Loading
Loading