1use std::sync::Arc;
19
20use arrow::{
21 array::{
22 Float64Builder, StringBuilder, TimestampNanosecondBuilder, UInt8Builder, UInt64Builder,
23 },
24 datatypes::{DataType, Field, Schema},
25 error::ArrowError,
26 record_batch::RecordBatch,
27};
28use nautilus_model::{data::OrderBookDelta, enums::BookAction};
29
30use super::{
31 float64_field, price_to_f64, quantity_to_f64, timestamp_field, unix_nanos_to_i64, utf8_field,
32};
33use crate::arrow::timestamp_data_type;
34
35#[must_use]
37pub fn deltas_schema() -> Schema {
38 Schema::new(vec![
39 utf8_field("instrument_id", false),
40 utf8_field("action", false),
41 utf8_field("side", false),
42 float64_field("price", false),
43 float64_field("size", false),
44 utf8_field("order_id", false),
45 Field::new("flags", DataType::UInt8, false),
46 Field::new("sequence", DataType::UInt64, false),
47 timestamp_field("ts_event", false),
48 timestamp_field("ts_init", false),
49 ])
50}
51
52pub fn encode_deltas(data: &[OrderBookDelta]) -> Result<RecordBatch, ArrowError> {
67 let mut instrument_id_builder = StringBuilder::new();
68 let mut action_builder = StringBuilder::new();
69 let mut side_builder = StringBuilder::new();
70 let mut price_builder = Float64Builder::with_capacity(data.len());
71 let mut size_builder = Float64Builder::with_capacity(data.len());
72 let mut order_id_builder = StringBuilder::new();
73 let mut flags_builder = UInt8Builder::with_capacity(data.len());
74 let mut sequence_builder = UInt64Builder::with_capacity(data.len());
75 let mut ts_event_builder =
76 TimestampNanosecondBuilder::with_capacity(data.len()).with_data_type(timestamp_data_type());
77 let mut ts_init_builder =
78 TimestampNanosecondBuilder::with_capacity(data.len()).with_data_type(timestamp_data_type());
79
80 for delta in data {
81 instrument_id_builder.append_value(delta.instrument_id.to_string());
82 action_builder.append_value(format!("{}", delta.action));
83 side_builder.append_value(
84 delta
85 .order
86 .side
87 .as_ref()
88 .map_or("NO_ORDER_SIDE", AsRef::as_ref),
89 );
90
91 if delta.action == BookAction::Clear {
95 price_builder.append_value(f64::NAN);
96 size_builder.append_value(f64::NAN);
97 } else {
98 price_builder.append_value(price_to_f64(&delta.order.price));
99 size_builder.append_value(quantity_to_f64(&delta.order.size));
100 }
101 order_id_builder.append_value(delta.order.order_id.to_string());
102 flags_builder.append_value(delta.flags);
103 sequence_builder.append_value(delta.sequence);
104 ts_event_builder.append_value(unix_nanos_to_i64(delta.ts_event.as_u64()));
105 ts_init_builder.append_value(unix_nanos_to_i64(delta.ts_init.as_u64()));
106 }
107
108 RecordBatch::try_new(
109 Arc::new(deltas_schema()),
110 vec![
111 Arc::new(instrument_id_builder.finish()),
112 Arc::new(action_builder.finish()),
113 Arc::new(side_builder.finish()),
114 Arc::new(price_builder.finish()),
115 Arc::new(size_builder.finish()),
116 Arc::new(order_id_builder.finish()),
117 Arc::new(flags_builder.finish()),
118 Arc::new(sequence_builder.finish()),
119 Arc::new(ts_event_builder.finish()),
120 Arc::new(ts_init_builder.finish()),
121 ],
122 )
123}
124
125#[cfg(test)]
126mod tests {
127 use arrow::{
128 array::{
129 Array, Float64Array, StringArray, TimestampNanosecondArray, UInt8Array, UInt64Array,
130 },
131 datatypes::TimeUnit,
132 };
133 use nautilus_model::{
134 data::order::BookOrder,
135 enums::{BookAction, OrderSide},
136 identifiers::InstrumentId,
137 types::{Price, Quantity, price::PRICE_UNDEF, quantity::QUANTITY_UNDEF},
138 };
139 use rstest::rstest;
140
141 use super::*;
142
143 fn make_delta(
144 instrument_id: &str,
145 action: BookAction,
146 side: OrderSide,
147 price: &str,
148 order_id: u64,
149 sequence: u64,
150 ts: u64,
151 ) -> OrderBookDelta {
152 OrderBookDelta {
153 instrument_id: InstrumentId::from(instrument_id),
154 action,
155 order: BookOrder {
156 side: Some(side),
157 price: Price::from(price),
158 size: Quantity::from(100),
159 order_id,
160 },
161 flags: 0,
162 sequence,
163 ts_event: ts.into(),
164 ts_init: (ts + 1).into(),
165 }
166 }
167
168 #[rstest]
169 fn test_encode_deltas_schema() {
170 let batch = encode_deltas(&[]).unwrap();
171 let fields = batch.schema().fields().clone();
172 assert_eq!(fields.len(), 10);
173 assert_eq!(fields[0].name(), "instrument_id");
174 assert_eq!(fields[0].data_type(), &DataType::Utf8);
175 assert_eq!(fields[1].name(), "action");
176 assert_eq!(fields[1].data_type(), &DataType::Utf8);
177 assert_eq!(fields[2].name(), "side");
178 assert_eq!(fields[2].data_type(), &DataType::Utf8);
179 assert_eq!(fields[3].name(), "price");
180 assert_eq!(fields[3].data_type(), &DataType::Float64);
181 assert_eq!(fields[4].name(), "size");
182 assert_eq!(fields[4].data_type(), &DataType::Float64);
183 assert_eq!(fields[5].name(), "order_id");
184 assert_eq!(fields[5].data_type(), &DataType::Utf8);
185 assert_eq!(fields[6].name(), "flags");
186 assert_eq!(fields[6].data_type(), &DataType::UInt8);
187 assert_eq!(fields[7].name(), "sequence");
188 assert_eq!(fields[7].data_type(), &DataType::UInt64);
189 assert_eq!(fields[8].name(), "ts_event");
190 assert_eq!(
191 fields[8].data_type(),
192 &DataType::Timestamp(TimeUnit::Nanosecond, Some("UTC".into()))
193 );
194 assert_eq!(fields[9].name(), "ts_init");
195 }
196
197 #[rstest]
198 fn test_encode_deltas_values() {
199 let deltas = vec![
200 make_delta(
201 "AAPL.XNAS",
202 BookAction::Add,
203 OrderSide::Buy,
204 "100.10",
205 1,
206 10,
207 1_000,
208 ),
209 make_delta(
210 "AAPL.XNAS",
211 BookAction::Update,
212 OrderSide::Sell,
213 "100.20",
214 2,
215 11,
216 2_000,
217 ),
218 ];
219 let batch = encode_deltas(&deltas).unwrap();
220
221 assert_eq!(batch.num_rows(), 2);
222
223 let action_col = batch
224 .column(1)
225 .as_any()
226 .downcast_ref::<StringArray>()
227 .unwrap();
228 let side_col = batch
229 .column(2)
230 .as_any()
231 .downcast_ref::<StringArray>()
232 .unwrap();
233 let price_col = batch
234 .column(3)
235 .as_any()
236 .downcast_ref::<Float64Array>()
237 .unwrap();
238 let size_col = batch
239 .column(4)
240 .as_any()
241 .downcast_ref::<Float64Array>()
242 .unwrap();
243 let order_id_col = batch
244 .column(5)
245 .as_any()
246 .downcast_ref::<StringArray>()
247 .unwrap();
248 let flags_col = batch
249 .column(6)
250 .as_any()
251 .downcast_ref::<UInt8Array>()
252 .unwrap();
253 let sequence_col = batch
254 .column(7)
255 .as_any()
256 .downcast_ref::<UInt64Array>()
257 .unwrap();
258 let ts_event_col = batch
259 .column(8)
260 .as_any()
261 .downcast_ref::<TimestampNanosecondArray>()
262 .unwrap();
263
264 assert_eq!(action_col.value(0), format!("{}", BookAction::Add));
265 assert_eq!(action_col.value(1), format!("{}", BookAction::Update));
266 assert_eq!(side_col.value(0), format!("{}", OrderSide::Buy));
267 assert_eq!(side_col.value(1), format!("{}", OrderSide::Sell));
268 assert!((price_col.value(0) - 100.10).abs() < 1e-9);
269 assert!((price_col.value(1) - 100.20).abs() < 1e-9);
270 assert!((size_col.value(0) - 100.0).abs() < 1e-9);
271 assert_eq!(order_id_col.value(0), "1");
272 assert_eq!(order_id_col.value(1), "2");
273 assert_eq!(flags_col.value(0), 0);
274 assert_eq!(sequence_col.value(0), 10);
275 assert_eq!(sequence_col.value(1), 11);
276 assert_eq!(ts_event_col.value(0), 1_000);
277 }
278
279 #[rstest]
280 fn test_encode_deltas_empty() {
281 let batch = encode_deltas(&[]).unwrap();
282 assert_eq!(batch.num_rows(), 0);
283 }
284
285 #[rstest]
286 fn test_encode_deltas_live_clear_renders_as_nan() {
287 let clear = OrderBookDelta::clear(InstrumentId::from("AAPL.XNAS"), 1, 1.into(), 2.into());
288
289 let batch = encode_deltas(&[clear]).unwrap();
290 let side_col = batch
291 .column(2)
292 .as_any()
293 .downcast_ref::<StringArray>()
294 .unwrap();
295 let price_col = batch
296 .column(3)
297 .as_any()
298 .downcast_ref::<Float64Array>()
299 .unwrap();
300 let size_col = batch
301 .column(4)
302 .as_any()
303 .downcast_ref::<Float64Array>()
304 .unwrap();
305
306 assert_eq!(side_col.value(0), "NO_ORDER_SIDE");
307 assert!(
308 price_col.value(0).is_nan(),
309 "live clear price should be NaN"
310 );
311 assert!(size_col.value(0).is_nan(), "live clear size should be NaN");
312 }
313
314 #[rstest]
315 fn test_encode_deltas_clear_sentinels_render_as_nan() {
316 let clear = OrderBookDelta {
317 instrument_id: InstrumentId::from("AAPL.XNAS"),
318 action: BookAction::Clear,
319 order: BookOrder {
320 side: None,
321 price: Price::from_raw(PRICE_UNDEF, 0),
322 size: Quantity::from_raw(QUANTITY_UNDEF, 0),
323 order_id: 0,
324 },
325 flags: 0,
326 sequence: 1,
327 ts_event: 1.into(),
328 ts_init: 2.into(),
329 };
330
331 let batch = encode_deltas(&[clear]).unwrap();
332 let price_col = batch
333 .column(3)
334 .as_any()
335 .downcast_ref::<Float64Array>()
336 .unwrap();
337 let size_col = batch
338 .column(4)
339 .as_any()
340 .downcast_ref::<Float64Array>()
341 .unwrap();
342
343 assert!(price_col.value(0).is_nan(), "clear price should be NaN");
344 assert!(size_col.value(0).is_nan(), "clear size should be NaN");
345 }
346
347 #[rstest]
348 fn test_encode_deltas_mixed_instruments() {
349 let deltas = vec![
350 make_delta(
351 "AAPL.XNAS",
352 BookAction::Add,
353 OrderSide::Buy,
354 "100.10",
355 1,
356 1,
357 1,
358 ),
359 make_delta(
360 "MSFT.XNAS",
361 BookAction::Add,
362 OrderSide::Sell,
363 "250.00",
364 2,
365 1,
366 2,
367 ),
368 ];
369 let batch = encode_deltas(&deltas).unwrap();
370 let instrument_id_col = batch
371 .column(0)
372 .as_any()
373 .downcast_ref::<StringArray>()
374 .unwrap();
375 assert_eq!(instrument_id_col.value(0), "AAPL.XNAS");
376 assert_eq!(instrument_id_col.value(1), "MSFT.XNAS");
377 }
378}