Skip to content

Commit 396ca91

Browse files
committed
Improve table writer signature flexibility
1 parent e476d3e commit 396ca91

16 files changed

Lines changed: 56 additions & 1593 deletions

File tree

rust/Cargo.toml

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -158,9 +158,9 @@ mmap = ["dep:libc"]
158158
# Adds parquet IO
159159
parquet = []
160160

161-
# CSV reader/writer/encoder/decoder. Pulls in itoa+ryu+memchr for
161+
# CSV reader/writer/encoder/decoder. Pulls in itoa+ryu+memchr for
162162
# fast formatter and quote scanner, atoi+fast-float2 for SIMD-accelerated
163-
# integer and float parsing on the decode path.
163+
# integer and float parsing on the decode path.
164164
csv = ["dep:ryu", "dep:memchr", "dep:fast-float2"]
165165

166166
# TCP transport

rust/examples/arrow/table_reader.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -128,7 +128,7 @@ async fn write_stream(tables: &[Table], path: &Path) -> Result<(), Box<dyn std::
128128
let schema: Vec<Field> = tables[0].schema().iter().map(|f| (**f).clone()).collect();
129129
let mut writer = TableStreamWriter::<Vec64<u8>>::new(schema, IPCMessageProtocol::Stream, None);
130130
for table in tables {
131-
writer.write(&table.clone().into())?;
131+
writer.write(table.clone())?;
132132
}
133133
writer.finish()?;
134134

rust/examples/arrow/table_stream_reader.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -101,7 +101,7 @@ async fn write_stream(
101101
let schema: Vec<Field> = tables[0].schema().iter().map(|f| (**f).clone()).collect();
102102
let mut writer = TableStreamWriter::<Vec64<u8>>::new(schema, protocol, None);
103103
for table in tables {
104-
writer.write(&table.clone().into())?;
104+
writer.write(table.clone())?;
105105
}
106106
writer.finish()?;
107107

rust/examples/arrow/table_stream_writer.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
3333
let mut writer = TableStreamWriter::<Vec64<u8>>::new(schema, IPCMessageProtocol::Stream, None);
3434

3535
for table in &tables {
36-
writer.write(&table.clone().into())?;
36+
writer.write(table.clone())?;
3737
}
3838
writer.finish()?;
3939

rust/src/lib.rs

Lines changed: 0 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -130,10 +130,6 @@ pub mod models {
130130

131131
/// TLV sink for simple type-length-value framing
132132
pub mod tlv_sink;
133-
134-
/// Live LBuffer-backed table sink for decoded records
135-
#[cfg(all(feature = "lbuffer", feature = "json"))]
136-
pub mod live_table_sink;
137133
}
138134

139135
/// Encoders for Arrow IPC, TLV, CSV, and optionally Parquet

rust/src/models/encoders/ipc/table_stream.rs

Lines changed: 8 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -294,7 +294,7 @@ mod tests {
294294
4,
295295
);
296296

297-
writer.write(&tbl.clone().into()).unwrap();
297+
writer.write(tbl.clone()).unwrap();
298298
writer.finish().unwrap();
299299

300300
let mut file = StdFile::create(&path).unwrap();
@@ -343,7 +343,7 @@ mod tests {
343343
4,
344344
);
345345

346-
writer.write(&tbl.clone().into()).unwrap();
346+
writer.write(tbl.clone()).unwrap();
347347
writer.finish().unwrap();
348348

349349
let mut file = StdFile::create(&path).unwrap();
@@ -392,7 +392,7 @@ mod tests {
392392
4,
393393
);
394394

395-
writer.write(&tbl.into()).unwrap();
395+
writer.write(tbl).unwrap();
396396
writer.finish().unwrap();
397397

398398
let mut buf = Vec::new();
@@ -437,7 +437,7 @@ mod tests {
437437
4,
438438
);
439439

440-
writer.write(&tbl.into()).unwrap();
440+
writer.write(tbl).unwrap();
441441
writer.finish().unwrap();
442442

443443
let mut buf = Vec::new();
@@ -482,7 +482,7 @@ mod tests {
482482
4,
483483
);
484484

485-
writer.write(&tbl.into()).unwrap();
485+
writer.write(tbl).unwrap();
486486
writer.finish().unwrap();
487487

488488
let mut buf = Vec::new();
@@ -527,7 +527,7 @@ mod tests {
527527
4,
528528
);
529529

530-
writer.write(&tbl.clone().into()).unwrap();
530+
writer.write(tbl.clone()).unwrap();
531531
writer.finish().unwrap();
532532

533533
// Write to temp file
@@ -591,7 +591,7 @@ mod tests {
591591
4,
592592
);
593593

594-
writer.write(&tbl.clone().into()).unwrap();
594+
writer.write(tbl.clone()).unwrap();
595595
writer.finish().unwrap();
596596

597597
let mut file = StdFile::create(&path).unwrap();
@@ -654,7 +654,7 @@ mod tests {
654654
4,
655655
);
656656

657-
writer.write(&tbl.clone().into()).unwrap();
657+
writer.write(tbl.clone()).unwrap();
658658
writer.finish().unwrap();
659659

660660
// Write to temp file

rust/src/models/readers/ipc/table.rs

Lines changed: 13 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -273,8 +273,8 @@ mod tests {
273273
let mut writer =
274274
TableStreamWriter::<Vec64<u8>>::new(schema.clone(), IPCMessageProtocol::Stream, None);
275275
register_dictionaries_for_table(&mut writer, &table);
276-
writer.write(&table.clone().into()).unwrap();
277-
writer.write(&table.clone().into()).unwrap();
276+
writer.write(table.clone()).unwrap();
277+
writer.write(table.clone()).unwrap();
278278
writer.finish().unwrap();
279279
let frames = writer.drain_all_frames();
280280

@@ -305,9 +305,9 @@ mod tests {
305305
TableStreamWriter::<Vec64<u8>>::new(schema.clone(), IPCMessageProtocol::Stream, None);
306306
register_dictionaries_for_table(&mut writer, &table);
307307
// three batches
308-
writer.write(&table.clone().into()).unwrap();
309-
writer.write(&table.clone().into()).unwrap();
310-
writer.write(&table.clone().into()).unwrap();
308+
writer.write(table.clone()).unwrap();
309+
writer.write(table.clone()).unwrap();
310+
writer.write(table.clone()).unwrap();
311311
writer.finish().unwrap();
312312
let frames = writer.drain_all_frames();
313313

@@ -332,8 +332,8 @@ mod tests {
332332
let mut writer =
333333
TableStreamWriter::<Vec64<u8>>::new(schema.clone(), IPCMessageProtocol::Stream, None);
334334
register_dictionaries_for_table(&mut writer, &table);
335-
writer.write(&table.clone().into()).unwrap();
336-
writer.write(&table.clone().into()).unwrap();
335+
writer.write(table.clone()).unwrap();
336+
writer.write(table.clone()).unwrap();
337337
writer.finish().unwrap();
338338
let frames = writer.drain_all_frames();
339339

@@ -363,8 +363,8 @@ mod tests {
363363
let mut writer =
364364
TableStreamWriter::<Vec64<u8>>::new(schema.clone(), IPCMessageProtocol::Stream, None);
365365
register_dictionaries_for_table(&mut writer, &table);
366-
writer.write(&table.clone().into()).unwrap();
367-
writer.write(&table.clone().into()).unwrap();
366+
writer.write(table.clone()).unwrap();
367+
writer.write(table.clone()).unwrap();
368368
writer.finish().unwrap();
369369
let frames = writer.drain_all_frames();
370370

@@ -392,7 +392,7 @@ mod tests {
392392
let mut writer =
393393
TableStreamWriter::<Vec64<u8>>::new(schema.clone(), IPCMessageProtocol::Stream, None);
394394
register_dictionaries_for_table(&mut writer, &table);
395-
writer.write(&table.clone().into()).unwrap();
395+
writer.write(table.clone()).unwrap();
396396
writer.finish().unwrap();
397397
let frames = writer.drain_all_frames();
398398

@@ -435,8 +435,8 @@ mod tests {
435435
let mut writer =
436436
TableStreamWriter::<Vec64<u8>>::new(schema.clone(), IPCMessageProtocol::Stream, None);
437437
register_dictionaries_for_table(&mut writer, &table);
438-
writer.write(&table.clone().into()).unwrap();
439-
writer.write(&table.clone().into()).unwrap();
438+
writer.write(table.clone()).unwrap();
439+
writer.write(table.clone()).unwrap();
440440
writer.finish().unwrap();
441441
let frames = writer.drain_all_frames();
442442

@@ -465,7 +465,7 @@ mod tests {
465465
let mut writer =
466466
TableStreamWriter::<Vec64<u8>>::new(schema.clone(), IPCMessageProtocol::Stream, None);
467467
register_dictionaries_for_table(&mut writer, &table);
468-
writer.write(&table.clone().into()).unwrap();
468+
writer.write(table.clone()).unwrap();
469469
writer.finish().unwrap();
470470
let frames = writer.drain_all_frames();
471471

0 commit comments

Comments
 (0)