From f573225820dae47367fc497da262430711a32f9c Mon Sep 17 00:00:00 2001 From: Benjamin Philip Date: Thu, 30 Jul 2026 20:48:50 +0530 Subject: [PATCH 1/9] add record_batch:body_length/1 --- src/arrow_ipc_record_batch.erl | 13 +++++++++++-- test/arrow_ipc_record_batch_SUITE.erl | 7 ++++++- 2 files changed, 17 insertions(+), 3 deletions(-) diff --git a/src/arrow_ipc_record_batch.erl b/src/arrow_ipc_record_batch.erl index 6d1cd77..4c59b41 100644 --- a/src/arrow_ipc_record_batch.erl +++ b/src/arrow_ipc_record_batch.erl @@ -44,7 +44,7 @@ comapatibility. You can find RecordBatches in the Arrow spec [here](https://arrow.apache.org/docs/format/Columnar.html#recordbatch-message). """. --export([from_erlang/1]). +-export([from_erlang/1, body_length/1]). -export_type([field_node/0, buffer/0, record_batch/0]). -include("arrow_ipc_record_batch.hrl"). @@ -55,7 +55,7 @@ You can find RecordBatches in the Arrow spec -type record_batch() :: #record_batch{}. -doc """ -Creates a RecordBatch given a list of arrays +Creates a RecordBatch given a list of arrays. """. -spec from_erlang(Arrays :: [arrow_array:array()]) -> RecordBatch :: record_batch(). from_erlang(Arrays) -> @@ -114,3 +114,12 @@ buffer_data(Buffer, CurOffset) -> #{offset => CurOffset, length => Buffer#buffer.length}, arrow_buffer:size(Buffer) + CurOffset }. + +-doc """ +Returns the body length of a Record Batch. +""". +-spec body_length(RecordBatch :: record_batch()) -> non_neg_integer(). +body_length(RecordBatch) -> + Buffers = RecordBatch#record_batch.buffers, + #{offset := Offset, length := Length} = lists:last(Buffers), + Offset + Length + (8 - (Length rem 8)). diff --git a/test/arrow_ipc_record_batch_SUITE.erl b/test/arrow_ipc_record_batch_SUITE.erl index 2195e4f..2b74fef 100644 --- a/test/arrow_ipc_record_batch_SUITE.erl +++ b/test/arrow_ipc_record_batch_SUITE.erl @@ -32,7 +32,8 @@ all() -> valid_length_on_from_erlang, valid_nodes_on_from_erlang, valid_buffers_on_from_erlang, - valid_compression_on_from_erlang + valid_compression_on_from_erlang, + body_length ]. valid_length_on_from_erlang(_Config) -> @@ -62,3 +63,7 @@ valid_buffers_on_from_erlang(_Config) -> valid_compression_on_from_erlang(_Config) -> ?assertEqual((?RecordBatch)#record_batch.compression, undefined). + +body_length(_Config) -> + erlang:display(?Body), + ?assertEqual(arrow_ipc_record_batch:body_length(?RecordBatch), byte_size(?Body)). From 42f5f4d564cbfed339284357686d5fd6bb7bd7c0 Mon Sep 17 00:00:00 2001 From: Benjamin Philip Date: Thu, 30 Jul 2026 20:50:09 +0530 Subject: [PATCH 2/9] stash --- src/arrow_ipc_message.erl | 11 ++++++----- src/arrow_ipc_message.hrl | 2 +- test/arrow_ipc_marks_data.hrl | 4 ++-- test/arrow_ipc_message_SUITE.erl | 2 +- 4 files changed, 10 insertions(+), 9 deletions(-) diff --git a/src/arrow_ipc_message.erl b/src/arrow_ipc_message.erl index c27eb14..c35155a 100644 --- a/src/arrow_ipc_message.erl +++ b/src/arrow_ipc_message.erl @@ -33,7 +33,7 @@ to represent a message. Metadata such as: 3. `body_length`: The length of the body in bytes 4. `custom_metadata`: A list of custom metadata in key-value format 5. `body`: The actual body. Can be undefined (in the case of Schema) - or a binary (in the case of Record Batch). + or a `t:arrow_array:array/0` (in the case of Record Batch). Currently, changing the version and custom metadata are not supported, but they have been added for forwards compatibility. @@ -79,10 +79,11 @@ from_erlang(Header) -> Creates a message given a data header and a body. """. -spec from_erlang( - Header :: arrow_ipc_schema:schema() | arrow_ipc_record_batch:record_batch(), Body :: binary() + Header :: arrow_ipc_schema:schema() | arrow_ipc_record_batch:record_batch(), + Body :: [arrow_array:array()] ) -> Message :: message(). from_erlang(Header, Body) -> - #message{header = Header, body = Body, body_length = byte_size(Body)}. + #message{header = Header, body = Body, body_length = undefined}. -doc """ Serializes a message into the Encapsulated Message Format. @@ -98,8 +99,8 @@ to_ipc(Message) -> case Message#message.body of undefined -> <<>>; - Bin -> - Bin + Arrays -> + body_from_erlang(Arrays) end, <>. diff --git a/src/arrow_ipc_message.hrl b/src/arrow_ipc_message.hrl index d818269..08bd187 100644 --- a/src/arrow_ipc_message.hrl +++ b/src/arrow_ipc_message.hrl @@ -27,5 +27,5 @@ %% This field is unique to arrow. %% The rest are from the flatbuffers definitions. - body :: binary() | undefined + body :: [arrow_array:array()] | undefined }). diff --git a/test/arrow_ipc_marks_data.hrl b/test/arrow_ipc_marks_data.hrl index 39669f7..cb6ac56 100644 --- a/test/arrow_ipc_marks_data.hrl +++ b/test/arrow_ipc_marks_data.hrl @@ -61,9 +61,9 @@ ). -define(Columns, [?ID, ?Name, ?Age, ?AnnualMarks]). -define(RecordBatch, arrow_ipc_record_batch:from_erlang(?Columns)). --define(Body, <<<<(arrow_array:to_arrow(Array))/binary>> || Array <- ?Columns>>). --define(RecordBatchMsg, arrow_ipc_message:from_erlang(?RecordBatch, ?Body)). +-define(RecordBatchMsg, arrow_ipc_message:from_erlang(?RecordBatch, ?Columns)). -define(RecordBatchEMF, arrow_ipc_message:to_ipc(?RecordBatchMsg)). +-define(Body, <<<<(arrow_array:to_arrow(Array))/binary>> || Array <- ?Columns>>). %%%%%%%%%%%%%%%% %% IPC Stream %% diff --git a/test/arrow_ipc_message_SUITE.erl b/test/arrow_ipc_message_SUITE.erl index 27c928e..6116804 100644 --- a/test/arrow_ipc_message_SUITE.erl +++ b/test/arrow_ipc_message_SUITE.erl @@ -68,7 +68,7 @@ valid_custom_metadata_on_from_erlang(_Config) -> valid_body_on_from_erlang(_Config) -> ?assertEqual((?SchemaMsg)#message.body, undefined), - ?assertEqual((?RecordBatchMsg)#message.body, ?Body). + ?assertEqual((?RecordBatchMsg)#message.body, ?Columns). %%%%%%%%%%%%%% %% to_ipc/1 %% From 55dedb857205a4db626ac19ca0ced06efc4d99c1 Mon Sep 17 00:00:00 2001 From: Benjamin Philip Date: Mon, 3 Aug 2026 10:35:30 +0530 Subject: [PATCH 3/9] add arrow_array:size/1 --- src/arrow_array.erl | 36 ++++++++++++++++++++++++++++-------- test/arrow_array_SUITE.erl | 31 ++++++++++++++++++++++++++++++- 2 files changed, 58 insertions(+), 9 deletions(-) diff --git a/src/arrow_array.erl b/src/arrow_array.erl index 56a89a4..3f30159 100644 --- a/src/arrow_array.erl +++ b/src/arrow_array.erl @@ -73,7 +73,8 @@ new arrays, and to serialize arrays to arrow exist. validity_bitmap/1, offsets/1, data/1, - to_arrow/1 + to_arrow/1, + size/1 ]). -include("arrow_array.hrl"). @@ -185,16 +186,35 @@ and not IPC. """. -spec to_arrow(Array :: array()) -> Arrow :: binary(). to_arrow(Array) -> - Validity = some(validity_bitmap(Array)), - Offsets = some(offsets(Array)), - Data = some(data(Array)), + Validity = some_to_arrow(validity_bitmap(Array)), + Offsets = some_to_arrow(offsets(Array)), + Data = some_to_arrow(data(Array)), <>. --spec some(Value :: array() | arrow_buffer:buffer() | undefined) -> Binary :: binary(). -some(undefined) -> +-spec some_to_arrow(Value :: array() | arrow_buffer:buffer() | undefined) -> Binary :: binary(). +some_to_arrow(undefined) -> <<>>; -some(Buffer) when is_record(Buffer, buffer) -> +some_to_arrow(Buffer) when is_record(Buffer, buffer) -> arrow_buffer:to_arrow(Buffer); -some(Array) when is_record(Array, array) -> +some_to_arrow(Array) when is_record(Array, array) -> to_arrow(Array). + +-doc """ +Returns the size of an array once serialized +""". +-spec size(Array :: array()) -> Size :: non_neg_integer(). +size(Array) -> + Validity = some_size(validity_bitmap(Array)), + Offsets = some_size(offsets(Array)), + Data = some_size(data(Array)), + + Validity + Offsets + Data. + +-spec some_size(Value :: array() | arrow_buffer:buffer() | undefined) -> Size :: non_neg_integer(). +some_size(undefined) -> + 0; +some_size(Buffer) when is_record(Buffer, buffer) -> + arrow_buffer:size(Buffer); +some_size(Array) when is_record(Array, array) -> + arrow_array:size(Array). diff --git a/test/arrow_array_SUITE.erl b/test/arrow_array_SUITE.erl index 8879764..f34c5e7 100644 --- a/test/arrow_array_SUITE.erl +++ b/test/arrow_array_SUITE.erl @@ -40,7 +40,8 @@ all() -> valid_data_on_access, %% serialization tests - valid_binary_on_to_arrow + valid_binary_on_to_arrow, + valid_size_on_size ]. %%%%%%%%%%%%%%%%%%%%%%%%%% @@ -99,6 +100,8 @@ valid_data_on_access(_Config) -> valid_binary_on_to_arrow(_Config) -> %% Primitive Fixed-Size Array %% Also offsets is not allocated on no offsets + Array0 = arrow_array:from_erlang(fixed_primitive, [], {s, 8}), + ?assertEqual(arrow_array:to_arrow(Array0), <<>>), Array1 = arrow_array:from_erlang(fixed_primitive, [1, 2, undefined, 3], {s, 8}), Validity1 = pad(<<0:1, 0:1, 0:1, 0:1, 1:1, 0:1, 1:1, 1:1>>), @@ -145,6 +148,32 @@ valid_binary_on_to_arrow(_Config) -> Bin4 = <>, ?assertEqual(arrow_array:to_arrow(Array4), Bin4). +valid_size_on_size(_Config) -> + %% Primitive Fixed-Size Array + Array0 = arrow_array:from_erlang(fixed_primitive, [], {s, 8}), + ?assertEqual(arrow_array:size(Array0), 0), + + Array1 = arrow_array:from_erlang(fixed_primitive, [1, 2, undefined, 3], {s, 8}), + ?assertEqual(arrow_array:size(Array1), byte_size(arrow_array:to_arrow(Array1))), + + %% Variable-Size Binary + Array2 = arrow_array:from_erlang( + variable_binary, [<<1>>, <<2, 3>>, undefined, <<4, 5, 6>>], {s, 8} + ), + ?assertEqual(arrow_array:size(Array2), byte_size(arrow_array:to_arrow(Array2))), + + %% Fixed-Sized List + Array3 = arrow_array:from_erlang( + fixed_list, [[1, 2, undefined, 3], [4, 5, 6, 7]], {s, 8} + ), + ?assertEqual(arrow_array:size(Array3), byte_size(arrow_array:to_arrow(Array3))), + + %% Variable-Size List + Array4 = arrow_array:from_erlang( + variable_list, [[1, 2, undefined], [3, 4], undefined, [5]], {s, 8} + ), + ?assertEqual(arrow_array:size(Array4), byte_size(arrow_array:to_arrow(Array4))). + %%%%%%%%%%% %% Utils %% %%%%%%%%%%% From b53e23772844ee1c675f9be56340401af643fadc Mon Sep 17 00:00:00 2001 From: Benjamin Philip Date: Mon, 3 Aug 2026 11:28:02 +0530 Subject: [PATCH 4/9] fix errors --- src/arrow_ipc_message.erl | 10 +++++----- src/arrow_ipc_record_batch.erl | 2 +- test/arrow_ipc_record_batch_SUITE.erl | 1 - 3 files changed, 6 insertions(+), 7 deletions(-) diff --git a/src/arrow_ipc_message.erl b/src/arrow_ipc_message.erl index c35155a..51e8601 100644 --- a/src/arrow_ipc_message.erl +++ b/src/arrow_ipc_message.erl @@ -68,22 +68,22 @@ for more info: -type key_value() :: #{key => string(), value => string()}. -doc """ -Creates a message given a data header. +Creates a message given a schema data header. """. --spec from_erlang(Header :: arrow_ipc_schema:schema() | arrow_ipc_record_batch:record_batch()) -> +-spec from_erlang(Header :: arrow_ipc_schema:schema()) -> Message :: message(). from_erlang(Header) -> #message{header = Header, body_length = 0}. -doc """ -Creates a message given a data header and a body. +Creates a message given a record batch data header and a body. """. -spec from_erlang( - Header :: arrow_ipc_schema:schema() | arrow_ipc_record_batch:record_batch(), + Header :: arrow_ipc_record_batch:record_batch(), Body :: [arrow_array:array()] ) -> Message :: message(). from_erlang(Header, Body) -> - #message{header = Header, body = Body, body_length = undefined}. + #message{header = Header, body = Body, body_length = arrow_ipc_record_batch:body_length(Header)}. -doc """ Serializes a message into the Encapsulated Message Format. diff --git a/src/arrow_ipc_record_batch.erl b/src/arrow_ipc_record_batch.erl index 4c59b41..76ef318 100644 --- a/src/arrow_ipc_record_batch.erl +++ b/src/arrow_ipc_record_batch.erl @@ -122,4 +122,4 @@ Returns the body length of a Record Batch. body_length(RecordBatch) -> Buffers = RecordBatch#record_batch.buffers, #{offset := Offset, length := Length} = lists:last(Buffers), - Offset + Length + (8 - (Length rem 8)). + Offset + Length + (64 - (Length rem 64)). diff --git a/test/arrow_ipc_record_batch_SUITE.erl b/test/arrow_ipc_record_batch_SUITE.erl index 2b74fef..f564f4a 100644 --- a/test/arrow_ipc_record_batch_SUITE.erl +++ b/test/arrow_ipc_record_batch_SUITE.erl @@ -65,5 +65,4 @@ valid_compression_on_from_erlang(_Config) -> ?assertEqual((?RecordBatch)#record_batch.compression, undefined). body_length(_Config) -> - erlang:display(?Body), ?assertEqual(arrow_ipc_record_batch:body_length(?RecordBatch), byte_size(?Body)). From 0f30f79fde36465718786947d78e6c5c20c400ee Mon Sep 17 00:00:00 2001 From: Benjamin Philip Date: Mon, 3 Aug 2026 11:32:08 +0530 Subject: [PATCH 5/9] Revert "add arrow_array:size/1" This reverts commit 55dedb857205a4db626ac19ca0ced06efc4d99c1. --- src/arrow_array.erl | 36 ++++++++---------------------------- test/arrow_array_SUITE.erl | 31 +------------------------------ 2 files changed, 9 insertions(+), 58 deletions(-) diff --git a/src/arrow_array.erl b/src/arrow_array.erl index 3f30159..56a89a4 100644 --- a/src/arrow_array.erl +++ b/src/arrow_array.erl @@ -73,8 +73,7 @@ new arrays, and to serialize arrays to arrow exist. validity_bitmap/1, offsets/1, data/1, - to_arrow/1, - size/1 + to_arrow/1 ]). -include("arrow_array.hrl"). @@ -186,35 +185,16 @@ and not IPC. """. -spec to_arrow(Array :: array()) -> Arrow :: binary(). to_arrow(Array) -> - Validity = some_to_arrow(validity_bitmap(Array)), - Offsets = some_to_arrow(offsets(Array)), - Data = some_to_arrow(data(Array)), + Validity = some(validity_bitmap(Array)), + Offsets = some(offsets(Array)), + Data = some(data(Array)), <>. --spec some_to_arrow(Value :: array() | arrow_buffer:buffer() | undefined) -> Binary :: binary(). -some_to_arrow(undefined) -> +-spec some(Value :: array() | arrow_buffer:buffer() | undefined) -> Binary :: binary(). +some(undefined) -> <<>>; -some_to_arrow(Buffer) when is_record(Buffer, buffer) -> +some(Buffer) when is_record(Buffer, buffer) -> arrow_buffer:to_arrow(Buffer); -some_to_arrow(Array) when is_record(Array, array) -> +some(Array) when is_record(Array, array) -> to_arrow(Array). - --doc """ -Returns the size of an array once serialized -""". --spec size(Array :: array()) -> Size :: non_neg_integer(). -size(Array) -> - Validity = some_size(validity_bitmap(Array)), - Offsets = some_size(offsets(Array)), - Data = some_size(data(Array)), - - Validity + Offsets + Data. - --spec some_size(Value :: array() | arrow_buffer:buffer() | undefined) -> Size :: non_neg_integer(). -some_size(undefined) -> - 0; -some_size(Buffer) when is_record(Buffer, buffer) -> - arrow_buffer:size(Buffer); -some_size(Array) when is_record(Array, array) -> - arrow_array:size(Array). diff --git a/test/arrow_array_SUITE.erl b/test/arrow_array_SUITE.erl index f34c5e7..8879764 100644 --- a/test/arrow_array_SUITE.erl +++ b/test/arrow_array_SUITE.erl @@ -40,8 +40,7 @@ all() -> valid_data_on_access, %% serialization tests - valid_binary_on_to_arrow, - valid_size_on_size + valid_binary_on_to_arrow ]. %%%%%%%%%%%%%%%%%%%%%%%%%% @@ -100,8 +99,6 @@ valid_data_on_access(_Config) -> valid_binary_on_to_arrow(_Config) -> %% Primitive Fixed-Size Array %% Also offsets is not allocated on no offsets - Array0 = arrow_array:from_erlang(fixed_primitive, [], {s, 8}), - ?assertEqual(arrow_array:to_arrow(Array0), <<>>), Array1 = arrow_array:from_erlang(fixed_primitive, [1, 2, undefined, 3], {s, 8}), Validity1 = pad(<<0:1, 0:1, 0:1, 0:1, 1:1, 0:1, 1:1, 1:1>>), @@ -148,32 +145,6 @@ valid_binary_on_to_arrow(_Config) -> Bin4 = <>, ?assertEqual(arrow_array:to_arrow(Array4), Bin4). -valid_size_on_size(_Config) -> - %% Primitive Fixed-Size Array - Array0 = arrow_array:from_erlang(fixed_primitive, [], {s, 8}), - ?assertEqual(arrow_array:size(Array0), 0), - - Array1 = arrow_array:from_erlang(fixed_primitive, [1, 2, undefined, 3], {s, 8}), - ?assertEqual(arrow_array:size(Array1), byte_size(arrow_array:to_arrow(Array1))), - - %% Variable-Size Binary - Array2 = arrow_array:from_erlang( - variable_binary, [<<1>>, <<2, 3>>, undefined, <<4, 5, 6>>], {s, 8} - ), - ?assertEqual(arrow_array:size(Array2), byte_size(arrow_array:to_arrow(Array2))), - - %% Fixed-Sized List - Array3 = arrow_array:from_erlang( - fixed_list, [[1, 2, undefined, 3], [4, 5, 6, 7]], {s, 8} - ), - ?assertEqual(arrow_array:size(Array3), byte_size(arrow_array:to_arrow(Array3))), - - %% Variable-Size List - Array4 = arrow_array:from_erlang( - variable_list, [[1, 2, undefined], [3, 4], undefined, [5]], {s, 8} - ), - ?assertEqual(arrow_array:size(Array4), byte_size(arrow_array:to_arrow(Array4))). - %%%%%%%%%%% %% Utils %% %%%%%%%%%%% From d823dff21d14da4e9ce162b7b3e3b7038fe1d0e7 Mon Sep 17 00:00:00 2001 From: Benjamin Philip Date: Mon, 3 Aug 2026 11:52:31 +0530 Subject: [PATCH 6/9] Use `pad_len` to calculate padding Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> --- src/arrow_ipc_record_batch.erl | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/src/arrow_ipc_record_batch.erl b/src/arrow_ipc_record_batch.erl index 76ef318..6e79919 100644 --- a/src/arrow_ipc_record_batch.erl +++ b/src/arrow_ipc_record_batch.erl @@ -120,6 +120,11 @@ Returns the body length of a Record Batch. """. -spec body_length(RecordBatch :: record_batch()) -> non_neg_integer(). body_length(RecordBatch) -> - Buffers = RecordBatch#record_batch.buffers, - #{offset := Offset, length := Length} = lists:last(Buffers), - Offset + Length + (64 - (Length rem 64)). + case RecordBatch#record_batch.buffers of + [] -> + 0; + Buffers -> + #{offset := Offset, length := Length} = lists:last(Buffers), + End = Offset + Length, + End + arrow_utils:pad_len(End) + end. From 214f478fccb3162e42f41a8133b6a610a60b0310 Mon Sep 17 00:00:00 2001 From: Benjamin Philip Date: Mon, 3 Aug 2026 11:56:09 +0530 Subject: [PATCH 7/9] Fix typespec --- src/arrow_ipc_message.erl | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/arrow_ipc_message.erl b/src/arrow_ipc_message.erl index 51e8601..ba0cc02 100644 --- a/src/arrow_ipc_message.erl +++ b/src/arrow_ipc_message.erl @@ -33,7 +33,7 @@ to represent a message. Metadata such as: 3. `body_length`: The length of the body in bytes 4. `custom_metadata`: A list of custom metadata in key-value format 5. `body`: The actual body. Can be undefined (in the case of Schema) - or a `t:arrow_array:array/0` (in the case of Record Batch). + or a `[t:arrow_array:array/0]` (in the case of Record Batch). Currently, changing the version and custom metadata are not supported, but they have been added for forwards compatibility. From 7a54f0d7119f608304269b96ee64a12aea4dd549 Mon Sep 17 00:00:00 2001 From: Benjamin Philip Date: Mon, 3 Aug 2026 11:56:48 +0530 Subject: [PATCH 8/9] format --- src/arrow_ipc_message.erl | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/src/arrow_ipc_message.erl b/src/arrow_ipc_message.erl index ba0cc02..95adb42 100644 --- a/src/arrow_ipc_message.erl +++ b/src/arrow_ipc_message.erl @@ -83,7 +83,9 @@ Creates a message given a record batch data header and a body. Body :: [arrow_array:array()] ) -> Message :: message(). from_erlang(Header, Body) -> - #message{header = Header, body = Body, body_length = arrow_ipc_record_batch:body_length(Header)}. + #message{ + header = Header, body = Body, body_length = arrow_ipc_record_batch:body_length(Header) + }. -doc """ Serializes a message into the Encapsulated Message Format. From 88ab59d20fb2b236117aa09c6266c30a710526d3 Mon Sep 17 00:00:00 2001 From: Benjamin Philip Date: Mon, 3 Aug 2026 12:19:37 +0530 Subject: [PATCH 9/9] Fix Livebook guide --- guides/quick-run-through.livemd | 22 ++++++++-------------- 1 file changed, 8 insertions(+), 14 deletions(-) diff --git a/guides/quick-run-through.livemd b/guides/quick-run-through.livemd index a3d6220..b781f97 100644 --- a/guides/quick-run-through.livemd +++ b/guides/quick-run-through.livemd @@ -1,5 +1,7 @@ -