1use std::sync::LazyLock;
19
20use bytes::{Buf, BufMut};
21use chrono::{Datelike, Timelike};
22use either::{Either, for_both};
23use enum_as_inner::EnumAsInner;
24use risingwave_pb::data::PbDatum;
25
26use crate::array::ArrayImpl;
27use crate::log::LogSuppressor;
28use crate::row::Row;
29use crate::types::*;
30
31pub mod error;
32use error::ValueEncodingError;
33
34use self::column_aware_row_encoding::ColumnAwareSerde;
35pub mod column_aware_row_encoding;
36
37pub use crate::row::RowDeserializer as BasicDeserializer;
38use crate::vector::{decode_vector_payload, encode_vector_payload};
39
40pub type Result<T> = std::result::Result<T, ValueEncodingError>;
41
42#[derive(EnumAsInner)]
44pub enum ValueRowSerdeKind {
45 Basic,
47 ColumnAware,
49}
50
51pub trait ValueRowSerializer: Clone {
53 fn serialize(&self, row: impl Row) -> Vec<u8>;
54}
55
56pub trait ValueRowDeserializer: Clone {
58 fn deserialize(&self, encoded_bytes: &[u8]) -> Result<Vec<Datum>>;
59}
60
61#[derive(Clone)]
63pub struct EitherSerde(pub Either<BasicSerde, ColumnAwareSerde>);
64
65impl From<BasicSerde> for EitherSerde {
66 fn from(value: BasicSerde) -> Self {
67 Self(Either::Left(value))
68 }
69}
70impl From<ColumnAwareSerde> for EitherSerde {
71 fn from(value: ColumnAwareSerde) -> Self {
72 Self(Either::Right(value))
73 }
74}
75
76impl ValueRowSerializer for EitherSerde {
77 fn serialize(&self, row: impl Row) -> Vec<u8> {
78 for_both!(&self.0, s => s.serialize(row))
79 }
80}
81
82impl ValueRowDeserializer for EitherSerde {
83 fn deserialize(&self, encoded_bytes: &[u8]) -> Result<Vec<Datum>> {
84 for_both!(&self.0, s => s.deserialize(encoded_bytes))
85 }
86}
87
88#[derive(Clone)]
90pub struct BasicSerializer;
91
92impl ValueRowSerializer for BasicSerializer {
93 fn serialize(&self, row: impl Row) -> Vec<u8> {
94 let mut buf = vec![];
95 for datum in row.iter() {
96 serialize_datum_into(datum, &mut buf);
97 }
98 buf
99 }
100}
101
102impl ValueRowDeserializer for BasicDeserializer {
103 fn deserialize(&self, encoded_bytes: &[u8]) -> Result<Vec<Datum>> {
104 Ok(self.deserialize(encoded_bytes)?.into_inner().into())
105 }
106}
107
108#[derive(Clone)]
110pub struct BasicSerde {
111 pub serializer: BasicSerializer,
112 pub deserializer: BasicDeserializer,
113}
114
115impl ValueRowSerializer for BasicSerde {
116 fn serialize(&self, row: impl Row) -> Vec<u8> {
117 self.serializer.serialize(row)
118 }
119}
120
121impl ValueRowDeserializer for BasicSerde {
122 fn deserialize(&self, encoded_bytes: &[u8]) -> Result<Vec<Datum>> {
123 Ok(self
124 .deserializer
125 .deserialize(encoded_bytes)?
126 .into_inner()
127 .into())
128 }
129}
130
131pub fn try_get_exact_serialize_datum_size(arr: &ArrayImpl) -> Option<usize> {
132 match arr {
133 ArrayImpl::Int16(_) => Some(2),
134 ArrayImpl::Int32(_) => Some(4),
135 ArrayImpl::Int64(_) => Some(8),
136 ArrayImpl::Serial(_) => Some(8),
137 ArrayImpl::Float32(_) => Some(4),
138 ArrayImpl::Float64(_) => Some(8),
139 ArrayImpl::Bool(_) => Some(1),
140 ArrayImpl::Decimal(_) => Some(estimate_serialize_decimal_size()),
141 ArrayImpl::Interval(_) => Some(estimate_serialize_interval_size()),
142 ArrayImpl::Date(_) => Some(estimate_serialize_date_size()),
143 ArrayImpl::Timestamp(_) => Some(estimate_serialize_timestamp_size()),
144 ArrayImpl::Time(_) => Some(estimate_serialize_time_size()),
145 _ => None,
146 }
147 .map(|x| x + 1)
148}
149
150pub fn serialize_datum(cell: impl ToDatumRef) -> Vec<u8> {
152 let mut buf: Vec<u8> = vec![];
153 serialize_datum_into(cell, &mut buf);
154 buf
155}
156
157pub fn serialize_datum_into(datum_ref: impl ToDatumRef, buf: &mut impl BufMut) {
159 if let Some(d) = datum_ref.to_datum_ref() {
160 buf.put_u8(1);
161 serialize_scalar(d, buf)
162 } else {
163 buf.put_u8(0);
164 }
165}
166
167pub fn estimate_serialize_datum_size(datum_ref: impl ToDatumRef) -> usize {
168 if let Some(d) = datum_ref.to_datum_ref() {
169 1 + estimate_serialize_scalar_size(d)
170 } else {
171 1
172 }
173}
174
175#[easy_ext::ext(DatumFromProtoExt)]
176impl Datum {
177 pub fn from_protobuf(proto: &PbDatum, data_type: &DataType) -> Result<Datum> {
179 deserialize_datum(proto.body.as_slice(), data_type)
180 }
181}
182
183#[easy_ext::ext(DatumToProtoExt)]
184impl<D: ToDatumRef> D {
185 pub fn to_protobuf(&self) -> PbDatum {
187 PbDatum {
188 body: serialize_datum(self),
189 }
190 }
191}
192
193pub fn deserialize_datum(mut data: impl Buf, ty: &DataType) -> Result<Datum> {
195 inner_deserialize_datum(&mut data, ty)
196}
197
198#[inline(always)]
200fn inner_deserialize_datum(data: &mut impl Buf, ty: &DataType) -> Result<Datum> {
201 let null_tag = data.get_u8();
202 match null_tag {
203 0 => Ok(None),
204 1 => Some(deserialize_value(ty, data)).transpose(),
205 _ => Err(ValueEncodingError::InvalidTagEncoding(null_tag)),
206 }
207}
208
209fn serialize_scalar(value: ScalarRefImpl<'_>, buf: &mut impl BufMut) {
210 match value {
211 ScalarRefImpl::Int16(v) => buf.put_i16_le(v),
212 ScalarRefImpl::Int32(v) => buf.put_i32_le(v),
213 ScalarRefImpl::Int64(v) => buf.put_i64_le(v),
214 ScalarRefImpl::Int256(v) => buf.put_slice(&v.to_le_bytes()),
215 ScalarRefImpl::Serial(v) => buf.put_i64_le(v.into_inner()),
216 ScalarRefImpl::Float32(v) => buf.put_f32_le(v.into_inner()),
217 ScalarRefImpl::Float64(v) => buf.put_f64_le(v.into_inner()),
218 ScalarRefImpl::Utf8(v) => serialize_str(v.as_bytes(), buf),
219 ScalarRefImpl::Bytea(v) => serialize_str(v, buf),
220 ScalarRefImpl::Bool(v) => buf.put_u8(v as u8),
221 ScalarRefImpl::Decimal(v) => serialize_decimal(&v, buf),
222 ScalarRefImpl::Interval(v) => serialize_interval(&v, buf),
223 ScalarRefImpl::Date(v) => serialize_date(v.0.num_days_from_ce(), buf),
224 ScalarRefImpl::Timestamp(v) => serialize_timestamp(
225 v.0.and_utc().timestamp(),
226 v.0.and_utc().timestamp_subsec_nanos(),
227 buf,
228 ),
229 ScalarRefImpl::Timestamptz(v) => buf.put_i64_le(v.timestamp_micros()),
230 ScalarRefImpl::Time(v) => {
231 serialize_time(v.0.num_seconds_from_midnight(), v.0.nanosecond(), buf)
232 }
233 ScalarRefImpl::Jsonb(v) => serialize_str(&v.value_serialize(), buf),
234 ScalarRefImpl::Variant(v) => serialize_str(v.as_bytes(), buf),
235 ScalarRefImpl::Struct(s) => serialize_struct(s, buf),
236 ScalarRefImpl::List(v) => serialize_list(v, buf),
237 ScalarRefImpl::Map(m) => serialize_list(m.into_inner(), buf),
238 ScalarRefImpl::Vector(v) => serialize_vector(v, buf),
239 }
240}
241
242fn estimate_serialize_scalar_size(value: ScalarRefImpl<'_>) -> usize {
243 match value {
244 ScalarRefImpl::Int16(_) => 2,
245 ScalarRefImpl::Int32(_) => 4,
246 ScalarRefImpl::Int64(_) => 8,
247 ScalarRefImpl::Int256(_) => 32,
248 ScalarRefImpl::Serial(_) => 8,
249 ScalarRefImpl::Float32(_) => 4,
250 ScalarRefImpl::Float64(_) => 8,
251 ScalarRefImpl::Utf8(v) => estimate_serialize_str_size(v.as_bytes()),
252 ScalarRefImpl::Bytea(v) => estimate_serialize_str_size(v),
253 ScalarRefImpl::Bool(_) => 1,
254 ScalarRefImpl::Decimal(_) => estimate_serialize_decimal_size(),
255 ScalarRefImpl::Interval(_) => estimate_serialize_interval_size(),
256 ScalarRefImpl::Date(_) => estimate_serialize_date_size(),
257 ScalarRefImpl::Timestamp(_) => estimate_serialize_timestamp_size(),
258 ScalarRefImpl::Timestamptz(_) => 8,
259 ScalarRefImpl::Time(_) => estimate_serialize_time_size(),
260 ScalarRefImpl::Jsonb(v) => v.capacity(),
262 ScalarRefImpl::Variant(v) => estimate_serialize_str_size(v.as_bytes()),
263 ScalarRefImpl::Struct(s) => estimate_serialize_struct_size(s),
264 ScalarRefImpl::List(v) => estimate_serialize_list_size(v),
265 ScalarRefImpl::Map(v) => estimate_serialize_list_size(v.into_inner()),
266 ScalarRefImpl::Vector(v) => estimate_serialize_vector_size(v),
267 }
268}
269
270fn serialize_struct(value: StructRef<'_>, buf: &mut impl BufMut) {
271 value.iter_fields_ref().for_each(|field_value| {
272 serialize_datum_into(field_value, buf);
273 });
274}
275
276fn estimate_serialize_struct_size(s: StructRef<'_>) -> usize {
277 s.estimate_serialize_size_inner()
278}
279fn serialize_list(value: ListRef<'_>, buf: &mut impl BufMut) {
280 let elems = value.iter();
281 buf.put_u32_le(elems.len() as u32);
282
283 elems.for_each(|field_value| {
284 serialize_datum_into(field_value, buf);
285 });
286}
287fn estimate_serialize_list_size(list: ListRef<'_>) -> usize {
288 4 + list.estimate_serialize_size_inner()
289}
290
291fn serialize_vector(value: VectorRef<'_>, buf: &mut impl BufMut) {
292 let elems = value.as_slice();
293 encode_vector_payload(elems, buf);
294}
295fn estimate_serialize_vector_size(v: VectorRef<'_>) -> usize {
296 size_of_val(v.as_slice())
297}
298
299fn serialize_str(bytes: &[u8], buf: &mut impl BufMut) {
300 buf.put_u32_le(bytes.len() as u32);
301 buf.put_slice(bytes);
302}
303
304fn estimate_serialize_str_size(bytes: &[u8]) -> usize {
305 4 + bytes.len()
306}
307
308fn serialize_interval(interval: &Interval, buf: &mut impl BufMut) {
309 buf.put_i32_le(interval.months());
310 buf.put_i32_le(interval.days());
311 buf.put_i64_le(interval.usecs());
312}
313
314fn estimate_serialize_interval_size() -> usize {
315 4 + 4 + 8
316}
317
318fn serialize_date(days: i32, buf: &mut impl BufMut) {
319 buf.put_i32_le(days);
320}
321
322fn estimate_serialize_date_size() -> usize {
323 4
324}
325
326fn serialize_timestamp(secs: i64, nsecs: u32, buf: &mut impl BufMut) {
327 buf.put_i64_le(secs);
328 buf.put_u32_le(nsecs);
329}
330
331fn estimate_serialize_timestamp_size() -> usize {
332 8 + 4
333}
334
335fn serialize_time(secs: u32, nano: u32, buf: &mut impl BufMut) {
336 buf.put_u32_le(secs);
337 buf.put_u32_le(nano);
338}
339
340fn estimate_serialize_time_size() -> usize {
341 4 + 4
342}
343
344fn serialize_decimal(decimal: &Decimal, buf: &mut impl BufMut) {
345 buf.put_slice(&decimal.unordered_serialize());
346}
347
348fn estimate_serialize_decimal_size() -> usize {
349 16
350}
351
352fn deserialize_value(ty: &DataType, data: &mut impl Buf) -> Result<ScalarImpl> {
353 Ok(match ty {
354 DataType::Int16 => ScalarImpl::Int16(data.get_i16_le()),
355 DataType::Int32 => ScalarImpl::Int32(data.get_i32_le()),
356 DataType::Int64 => ScalarImpl::Int64(data.get_i64_le()),
357 DataType::Int256 => ScalarImpl::Int256(deserialize_int256(data)),
358 DataType::Serial => ScalarImpl::Serial(Serial::from(data.get_i64_le())),
359 DataType::Float32 => ScalarImpl::Float32(F32::from(data.get_f32_le())),
360 DataType::Float64 => ScalarImpl::Float64(F64::from(data.get_f64_le())),
361 DataType::Varchar => ScalarImpl::Utf8(deserialize_str(data)?),
362 DataType::Boolean => ScalarImpl::Bool(deserialize_bool(data)?),
363 DataType::Decimal => ScalarImpl::Decimal(deserialize_decimal(data)?),
364 DataType::Interval => ScalarImpl::Interval(deserialize_interval(data)?),
365 DataType::Time => ScalarImpl::Time(deserialize_time(data)?),
366 DataType::Timestamp => ScalarImpl::Timestamp(deserialize_timestamp(data)?),
367 DataType::Timestamptz => {
368 let micros = data.get_i64_le();
369 let tz = Timestamptz::from_micros(micros).unwrap_or_else(|| {
370 static LOG_SUPPRESSOR: LazyLock<LogSuppressor> =
373 LazyLock::new(LogSuppressor::default);
374 if let Ok(suppressed_count) = LOG_SUPPRESSOR.check() {
375 tracing::error!(suppressed_count, micros, "decoded out-of-range timestamptz");
376 }
377 Timestamptz::from_micros_uncheck(micros)
378 });
379 ScalarImpl::Timestamptz(tz)
380 }
381 DataType::Date => ScalarImpl::Date(deserialize_date(data)?),
382 DataType::Jsonb => ScalarImpl::Jsonb(
383 JsonbVal::value_deserialize(&deserialize_bytea(data))
384 .ok_or(ValueEncodingError::InvalidJsonbEncoding)?,
385 ),
386 DataType::Variant => ScalarImpl::Variant(
387 VariantVal::value_deserialize(&deserialize_bytea(data))
388 .ok_or(ValueEncodingError::InvalidVariantEncoding)?,
389 ),
390 DataType::Struct(struct_def) => deserialize_struct(struct_def, data)?,
391 DataType::Bytea => ScalarImpl::Bytea(deserialize_bytea(data).into()),
392 DataType::Vector(dimension) => deserialize_vector(*dimension, data),
393 DataType::List(list_type) => deserialize_list(list_type, data)?,
394 DataType::Map(map_type) => deserialize_map(map_type, data)?,
395 })
396}
397
398fn deserialize_struct(struct_def: &StructType, data: &mut impl Buf) -> Result<ScalarImpl> {
399 let mut field_values = Vec::with_capacity(struct_def.len());
400 for field_type in struct_def.types() {
401 field_values.push(inner_deserialize_datum(data, field_type)?);
402 }
403
404 Ok(ScalarImpl::Struct(StructValue::new(field_values)))
405}
406
407fn deserialize_list(list_type: &ListType, data: &mut impl Buf) -> Result<ScalarImpl> {
408 let elem_type = list_type.elem();
409 let len = data.get_u32_le();
410 let mut builder = elem_type.create_array_builder(len as usize);
411 for _ in 0..len {
412 builder.append(inner_deserialize_datum(data, elem_type)?);
413 }
414 Ok(ScalarImpl::List(ListValue::new(builder.finish())))
415}
416
417fn deserialize_map(map_type: &MapType, data: &mut impl Buf) -> Result<ScalarImpl> {
418 let list = deserialize_list(&map_type.clone().into_list_type(), data)?.into_list();
420 Ok(ScalarImpl::Map(MapValue::from_entries(list)))
421}
422
423fn deserialize_vector(dimension: usize, data: &mut impl Buf) -> ScalarImpl {
424 VectorVal {
425 inner: decode_vector_payload(dimension, data).into_boxed_slice(),
426 }
427 .to_scalar_value()
428}
429
430fn deserialize_str(data: &mut impl Buf) -> Result<Box<str>> {
431 let len = data.get_u32_le();
432 let mut bytes = vec![0; len as usize];
433 data.copy_to_slice(&mut bytes);
434 String::from_utf8(bytes)
435 .map(String::into_boxed_str)
436 .map_err(ValueEncodingError::InvalidUtf8)
437}
438
439fn deserialize_bytea(data: &mut impl Buf) -> Vec<u8> {
440 let len = data.get_u32_le();
441 let mut bytes = vec![0; len as usize];
442 data.copy_to_slice(&mut bytes);
443 bytes
444}
445
446fn deserialize_int256(data: &mut impl Buf) -> Int256 {
447 let mut bytes = [0; Int256::size()];
448 data.copy_to_slice(&mut bytes);
449 Int256::from_le_bytes(bytes)
450}
451
452fn deserialize_bool(data: &mut impl Buf) -> Result<bool> {
453 match data.get_u8() {
454 1 => Ok(true),
455 0 => Ok(false),
456 value => Err(ValueEncodingError::InvalidBoolEncoding(value)),
457 }
458}
459
460fn deserialize_interval(data: &mut impl Buf) -> Result<Interval> {
461 let months = data.get_i32_le();
462 let days = data.get_i32_le();
463 let usecs = data.get_i64_le();
464 Ok(Interval::from_month_day_usec(months, days, usecs))
465}
466
467fn deserialize_time(data: &mut impl Buf) -> Result<Time> {
468 let secs = data.get_u32_le();
469 let nano = data.get_u32_le();
470 Time::with_secs_nano(secs, nano)
471 .map_err(|_e| ValueEncodingError::InvalidTimeEncoding(secs, nano))
472}
473
474fn deserialize_timestamp(data: &mut impl Buf) -> Result<Timestamp> {
475 let secs = data.get_i64_le();
476 let nsecs = data.get_u32_le();
477 Timestamp::with_secs_nsecs(secs, nsecs)
478 .map_err(|_e| ValueEncodingError::InvalidTimestampEncoding(secs, nsecs))
479}
480
481fn deserialize_date(data: &mut impl Buf) -> Result<Date> {
482 let days = data.get_i32_le();
483 Date::with_days_since_ce(days).map_err(|_e| ValueEncodingError::InvalidDateEncoding(days))
484}
485
486fn deserialize_decimal(data: &mut impl Buf) -> Result<Decimal> {
487 let mut bytes = [0; 16];
488 data.copy_to_slice(&mut bytes);
489 Ok(Decimal::unordered_deserialize(bytes))
490}
491
492#[cfg(test)]
493mod tests {
494 use crate::array::{ArrayImpl, ListValue, StructValue};
495 use crate::test_utils::rand_chunk;
496 use crate::types::{
497 DataType, Date, Datum, Decimal, Interval, ScalarImpl, Serial, Time, Timestamp,
498 };
499 use crate::util::value_encoding::{
500 estimate_serialize_datum_size, serialize_datum, try_get_exact_serialize_datum_size,
501 };
502
503 fn test_estimate_serialize_scalar_size(s: ScalarImpl) {
504 let d = Datum::from(s);
505 assert_eq!(estimate_serialize_datum_size(&d), serialize_datum(&d).len());
506 }
507
508 fn test_try_get_exact_serialize_datum_size(s: &ArrayImpl) {
509 let d = s.to_datum();
510 if let Some(ret) = try_get_exact_serialize_datum_size(s) {
511 assert_eq!(ret, serialize_datum(&d).len());
512 }
513 }
514
515 #[test]
516 fn test_estimate_size() {
517 let d: Datum = None;
518 assert_eq!(estimate_serialize_datum_size(&d), serialize_datum(&d).len());
519
520 test_estimate_serialize_scalar_size(ScalarImpl::Bool(true));
521 test_estimate_serialize_scalar_size(ScalarImpl::Int16(1));
522 test_estimate_serialize_scalar_size(ScalarImpl::Int32(1));
523 test_estimate_serialize_scalar_size(ScalarImpl::Int64(1));
524 test_estimate_serialize_scalar_size(ScalarImpl::Float32(1.0.into()));
525 test_estimate_serialize_scalar_size(ScalarImpl::Float64(1.0.into()));
526 test_estimate_serialize_scalar_size(ScalarImpl::Serial(Serial::from(i64::MIN)));
527
528 test_estimate_serialize_scalar_size(ScalarImpl::Utf8("abc".into()));
529 test_estimate_serialize_scalar_size(ScalarImpl::Utf8("".into()));
530 test_estimate_serialize_scalar_size(ScalarImpl::Decimal(Decimal::NegativeInf));
531 test_estimate_serialize_scalar_size(ScalarImpl::Decimal(Decimal::PositiveInf));
532 test_estimate_serialize_scalar_size(ScalarImpl::Decimal(Decimal::NaN));
533 test_estimate_serialize_scalar_size(ScalarImpl::Decimal(123123.into()));
534 test_estimate_serialize_scalar_size(ScalarImpl::Interval(Interval::from_month_day_usec(
535 7, 8, 9,
536 )));
537 test_estimate_serialize_scalar_size(ScalarImpl::Date(Date::from_ymd_uncheck(2333, 3, 3)));
538 test_estimate_serialize_scalar_size(ScalarImpl::Bytea("\\x233".as_bytes().into()));
539 test_estimate_serialize_scalar_size(ScalarImpl::Time(Time::from_hms_uncheck(2, 3, 3)));
540 test_estimate_serialize_scalar_size(ScalarImpl::Timestamp(
541 Timestamp::from_timestamp_uncheck(23333333, 2333),
542 ));
543 test_estimate_serialize_scalar_size(ScalarImpl::Interval(Interval::from_month_day_usec(
544 2, 3, 3333,
545 )));
546 test_estimate_serialize_scalar_size(ScalarImpl::Struct(StructValue::new(vec![
547 ScalarImpl::Int64(233).into(),
548 ScalarImpl::Float64(23.33.into()).into(),
549 ])));
550 test_estimate_serialize_scalar_size(ScalarImpl::List(ListValue::from_iter([233i64, 2333])));
551 }
552
553 #[test]
554 fn test_try_estimate_size() {
555 let chunk = rand_chunk::gen_chunk(
556 &[
557 DataType::Int16,
558 DataType::Int32,
559 DataType::Int64,
560 DataType::Serial,
561 DataType::Float32,
562 DataType::Float64,
563 DataType::Boolean,
564 DataType::Decimal,
565 DataType::Interval,
566 DataType::Time,
567 DataType::Timestamp,
568 DataType::Date,
569 ],
570 1,
571 0,
572 0.0,
573 );
574 for column in chunk.columns() {
575 test_try_get_exact_serialize_datum_size(column);
576 }
577 }
578}