Skip to content
Open
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
84 changes: 52 additions & 32 deletions datafusion/common/src/scalar/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2861,17 +2861,14 @@ impl ScalarValue {
($ARRAY_TY:ident, $SCALAR_TY:ident, $TZ:expr) => {{
{
let array = scalars
.map(|sv| {
if let ScalarValue::$SCALAR_TY(v, _) = sv {
Ok(v)
} else {
_exec_err!(
"Inconsistent types in ScalarValue::iter_to_array. \
Expected {:?}, got {:?}",
data_type,
sv
)
}
.map(|sv| match sv {
ScalarValue::$SCALAR_TY(v, tz) if &tz == $TZ => Ok(v),
sv => _exec_err!(
"Inconsistent types in ScalarValue::iter_to_array. \
Expected {:?}, got {:?}",
data_type,
sv
),
})
.collect::<Result<$ARRAY_TY>>()?;
Arc::new(array.with_timezone_opt($TZ.clone()))
Expand Down Expand Up @@ -3161,15 +3158,16 @@ impl ScalarValue {
}
DataType::FixedSizeBinary(size) => {
let array = scalars
.map(|sv| {
if let ScalarValue::FixedSizeBinary(_, v) = sv {
.map(|sv| match sv {
ScalarValue::FixedSizeBinary(inner_size, v)
if inner_size == *size =>
{
Ok(v)
} else {
_exec_err!(
"Inconsistent types in ScalarValue::iter_to_array. \
Expected {data_type}, got {sv:?}"
)
}
sv => _exec_err!(
"Inconsistent types in ScalarValue::iter_to_array. \
Expected {data_type}, got {sv:?}"
),
})
.collect::<Result<Vec<_>>>()?;
let array = FixedSizeBinaryArray::try_from_sparse_iter_with_size(
Expand Down Expand Up @@ -3220,10 +3218,16 @@ impl ScalarValue {
let array = scalars
.into_iter()
.map(|element: ScalarValue| match element {
ScalarValue::Decimal32(v1, _, _) => Ok(v1),
s => {
_internal_err!("Expected ScalarValue::Null element. Received {s:?}")
ScalarValue::Decimal32(value, inner_precision, inner_scale)
if inner_precision == precision && inner_scale == scale =>
{
Ok(value)
}
scalar => _exec_err!(
"Inconsistent types in ScalarValue::iter_to_array. Expected {:?}, got {:?}",
DataType::Decimal32(precision, scale),
scalar
),
})
.collect::<Result<Decimal32Array>>()?
.with_precision_and_scale(precision, scale)?;
Expand All @@ -3238,10 +3242,16 @@ impl ScalarValue {
let array = scalars
.into_iter()
.map(|element: ScalarValue| match element {
ScalarValue::Decimal64(v1, _, _) => Ok(v1),
s => {
_internal_err!("Expected ScalarValue::Null element. Received {s:?}")
ScalarValue::Decimal64(value, inner_precision, inner_scale)
if inner_precision == precision && inner_scale == scale =>
{
Ok(value)
}
scalar => _exec_err!(
"Inconsistent types in ScalarValue::iter_to_array. Expected {:?}, got {:?}",
DataType::Decimal64(precision, scale),
scalar
),
})
.collect::<Result<Decimal64Array>>()?
.with_precision_and_scale(precision, scale)?;
Expand All @@ -3256,10 +3266,16 @@ impl ScalarValue {
let array = scalars
.into_iter()
.map(|element: ScalarValue| match element {
ScalarValue::Decimal128(v1, _, _) => Ok(v1),
s => {
_internal_err!("Expected ScalarValue::Null element. Received {s:?}")
ScalarValue::Decimal128(value, inner_precision, inner_scale)
if inner_precision == precision && inner_scale == scale =>
{
Ok(value)
}
scalar => _exec_err!(
"Inconsistent types in ScalarValue::iter_to_array. Expected {:?}, got {:?}",
DataType::Decimal128(precision, scale),
scalar
),
})
.collect::<Result<Decimal128Array>>()?
.with_precision_and_scale(precision, scale)?;
Expand All @@ -3274,12 +3290,16 @@ impl ScalarValue {
let array = scalars
.into_iter()
.map(|element: ScalarValue| match element {
ScalarValue::Decimal256(v1, _, _) => Ok(v1),
s => {
_internal_err!(
"Expected ScalarValue::Decimal256 element. Received {s:?}"
)
ScalarValue::Decimal256(value, inner_precision, inner_scale)
if inner_precision == precision && inner_scale == scale =>
{
Ok(value)
}
scalar => _exec_err!(
"Inconsistent types in ScalarValue::iter_to_array. Expected {:?}, got {:?}",
DataType::Decimal256(precision, scale),
scalar
),
})
.collect::<Result<Decimal256Array>>()?
.with_precision_and_scale(precision, scale)?;
Expand Down
70 changes: 56 additions & 14 deletions datafusion/core/tests/fuzz_cases/record_batch_generator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -61,10 +61,13 @@ pub fn get_supported_types_columns(rng_seed: u64) -> Vec<ColumnDescr> {
ColumnDescr::new("time32_ms", DataType::Time32(TimeUnit::Millisecond)),
ColumnDescr::new("time64_us", DataType::Time64(TimeUnit::Microsecond)),
ColumnDescr::new("time64_ns", DataType::Time64(TimeUnit::Nanosecond)),
ColumnDescr::new("timestamp_s", DataType::Timestamp(TimeUnit::Second, None)),
ColumnDescr::new(
"timestamp_s",
DataType::Timestamp(TimeUnit::Second, Some(Arc::from("UTC"))),
),
ColumnDescr::new(
"timestamp_ms",
DataType::Timestamp(TimeUnit::Millisecond, None),
DataType::Timestamp(TimeUnit::Millisecond, Some(Arc::from("+05:30"))),
),
ColumnDescr::new(
"timestamp_us",
Expand Down Expand Up @@ -142,6 +145,9 @@ pub fn get_supported_types_columns(rng_seed: u64) -> Vec<ColumnDescr> {
ColumnDescr::new("binary", DataType::Binary),
ColumnDescr::new("large_binary", DataType::LargeBinary),
ColumnDescr::new("binaryview", DataType::BinaryView),
ColumnDescr::new("fixed_binary_1", DataType::FixedSizeBinary(1)),
ColumnDescr::new("fixed_binary_8", DataType::FixedSizeBinary(8)),
ColumnDescr::new("fixed_binary_32", DataType::FixedSizeBinary(32)),
ColumnDescr::new(
"dictionary_utf8_low",
DataType::Dictionary(Box::new(DataType::UInt64), Box::new(DataType::Utf8)),
Expand Down Expand Up @@ -249,6 +255,27 @@ macro_rules! generate_primitive_array {
}};
}

macro_rules! generate_timestamp_array {
($SELF:ident, $NUM_ROWS:ident, $MAX_NUM_DISTINCT:expr, $NULL_PCT:ident, $BATCH_GEN_RNG:ident, $ARRAY_GEN_RNG:ident, $ARROW_TYPE:ident, $TIMEZONE:ident) => {{
let array = generate_primitive_array!(
$SELF,
$NUM_ROWS,
$MAX_NUM_DISTINCT,
$NULL_PCT,
$BATCH_GEN_RNG,
$ARRAY_GEN_RNG,
$ARROW_TYPE
);
let array = array
.as_any()
.downcast_ref::<PrimitiveArray<$ARROW_TYPE>>()
.unwrap()
.clone()
.with_timezone_opt($TIMEZONE.clone());
Arc::new(array) as ArrayRef
}};
}

macro_rules! generate_dict {
($SELF:ident, $NUM_ROWS:ident, $MAX_NUM_DISTINCT:expr, $NULL_PCT:ident, $BATCH_GEN_RNG:ident, $ARRAY_GEN_RNG:ident, $ARROW_TYPE:ident, $VALUES: ident) => {{
debug_assert_eq!($VALUES.len(), $MAX_NUM_DISTINCT);
Expand Down Expand Up @@ -617,48 +644,52 @@ impl RecordBatchGenerator {
DurationNanosecondType
)
}
DataType::Timestamp(TimeUnit::Second, None) => {
generate_primitive_array!(
DataType::Timestamp(TimeUnit::Second, ref timezone) => {
generate_timestamp_array!(
self,
num_rows,
max_num_distinct,
null_pct,
batch_gen_rng,
array_gen_rng,
TimestampSecondType
TimestampSecondType,
timezone
)
}
DataType::Timestamp(TimeUnit::Millisecond, None) => {
generate_primitive_array!(
DataType::Timestamp(TimeUnit::Millisecond, ref timezone) => {
generate_timestamp_array!(
self,
num_rows,
max_num_distinct,
null_pct,
batch_gen_rng,
array_gen_rng,
TimestampMillisecondType
TimestampMillisecondType,
timezone
)
}
DataType::Timestamp(TimeUnit::Microsecond, None) => {
generate_primitive_array!(
DataType::Timestamp(TimeUnit::Microsecond, ref timezone) => {
generate_timestamp_array!(
self,
num_rows,
max_num_distinct,
null_pct,
batch_gen_rng,
array_gen_rng,
TimestampMicrosecondType
TimestampMicrosecondType,
timezone
)
}
DataType::Timestamp(TimeUnit::Nanosecond, None) => {
generate_primitive_array!(
DataType::Timestamp(TimeUnit::Nanosecond, ref timezone) => {
generate_timestamp_array!(
self,
num_rows,
max_num_distinct,
null_pct,
batch_gen_rng,
array_gen_rng,
TimestampNanosecondType
TimestampNanosecondType,
timezone
)
}
DataType::Utf8 | DataType::LargeUtf8 | DataType::Utf8View => {
Expand Down Expand Up @@ -697,6 +728,17 @@ impl RecordBatchGenerator {
_ => unreachable!(),
}
}
DataType::FixedSizeBinary(width) => {
let mut generator = BinaryArrayGenerator {
max_len: usize::try_from(width)
.expect("fixed-size binary width must be nonnegative"),
num_binaries: num_rows,
num_distinct_binaries: max_num_distinct,
null_pct,
rng: array_gen_rng,
};
generator.gen_fixed_size_binary()
}
DataType::Decimal32(precision, scale) => {
generate_decimal_array!(
self,
Expand Down
Loading