[](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-1)import time [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-2)import numpy as np [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-3)import pandas as pd [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-4)import pyarrow as pa [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-5) [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-6)def generate_data(total_size, ncols): [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-7) nrows = total_size / ncols / np.dtype('float64').itemsize [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-8) return pd.DataFrame({ [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-9) 'c' + str(i): np.random.randn(nrows) [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-10) for i in range(ncols) [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-11) }) [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-12) [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-13)KILOBYTE = 1 << 10 [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-14)MEGABYTE = KILOBYTE * KILOBYTE [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-15)DATA_SIZE = 1024 * MEGABYTE [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-16)NCOLS = 16 [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-17) [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-18)def get_timing(f, niter): [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-19) start = time.clock_gettime(time.CLOCK_REALTIME) [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-20) for i in range(niter): [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-21) f() [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-22) return (time.clock_gettime(time.CLOCK_REALTIME) - start) / NITER [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-23) [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-24)def read_as_dataframe(klass, source): [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-25) reader = klass(source) [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-26) table = reader.read_all() [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-27) return table.to_pandas() [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-28)NITER = 5 [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-29)results = [] [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-30) [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-31)CHUNKSIZES = [16 * KILOBYTE, 64 * KILOBYTE, 256 * KILOBYTE, MEGABYTE, 16 * MEGABYTE] [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-32) [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-33)for chunksize in CHUNKSIZES: [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-34) nchunks = DATA_SIZE // chunksize [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-35) batch = pa.RecordBatch.from_pandas(generate_data(chunksize, NCOLS)) [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-36) [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-37) sink = pa.InMemoryOutputStream() [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-38) stream_writer = pa.StreamWriter(sink, batch.schema) [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-39) [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-40) for i in range(nchunks): [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-41) stream_writer.write_batch(batch) [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-42) [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-43) source = sink.get_result() [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-44) [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-45) elapsed = get_timing(lambda: read_as_dataframe(pa.StreamReader, source), NITER) [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-46) [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-47) result = (chunksize, elapsed) [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-48) print(result) [](https://wesmckinney.com/blog/arrow-streaming-columnar/#cb9-49) results.append(result)