apache / apache/arrow-julia

`Arrow.Stream` not working with non-seekable I/O e.g. FIFOs and sockets

Open
#580 0 comments 0 reactions 0 assignees View on GitHub
Dominant language
Julia
Stars
312
Forks
78
PR merge metrics
No merged PRs in 30d

Description

I found that `Arrow.Stream` does not work with non-seekable I/O, which should be supported for streaming. Here are some MWEs.

### Named pipes

```bash
mkfifo /tmp/arrow_pipe

# Producer
julia -e '
using Arrow
open("/tmp/arrow_pipe", "w") do io
Arrow.write(io, (i = collect(1:10),); file=false) # streaming
end
' &

# Consumer
julia -e '
using Arrow
open("/tmp/arrow_pipe", "r") do io
for batch in Arrow.Stream(io)
println(batch)
end
end
'
# Result: no output, no error

rm /tmp/arrow_pipe
```

### Sockets

```julia
using Arrow, Sockets

server = listen(9999)
@async begin
conn = accept(server)
Arrow.write(conn, (i = collect(1:10),); file=false)
close(conn)
end

sock = connect(9999)

# This block hangs indefinitely. Press Ctrl-C to proceed.
for batch in Arrow.Stream(sock)
println(batch)
end

wait(t)
# ERROR: TaskFailedException
# nested task error: MethodError: no method matching position(::TCPSocket)
```

### Unix domain sockets

```julia
using Arrow, Sockets

server = listen("/tmp/arrow.sock")
@async begin
conn = accept(server)
Arrow.write(conn, (i = collect(1:10),); file=false)
close(conn)
end

sock = connect("/tmp/arrow.sock")

# This block hangs indefinitely. Press Ctrl-C to proceed.
for batch in Arrow.Stream(sock)
println(batch)
end

wait(t)
# ERROR: TaskFailedException
# nested task error: MethodError: no method matching position(::Base.PipeEndpoint)
```

The above demo should be reproducible with Arrow v2.8.0.

## Diagnosis

The first issue is caused `tobytes(io::IOStream)` unconditionally using `Mmap.mmap(io)`. For named pipes, `filesize(io)` returns `0`, so `Mmap.mmap(io)` silently returns an empty `UInt8[]`.

The rest is caused by `Base.write(io, msg, ...)` calling `position(io)` to record block positions for the file format footer, even when writing streaming format (`file=false`).

I will submit a PR soon.

Contributor guide

No contributing guide indexed for this repository

Research direction

Reproduce the named-pipe and socket MWEs, then trace Arrow.Stream through tobytes(io::IOStream) and Base.write(io, msg, ...) to see how Mmap.mmap and position(io) are used. Done means Arrow.Stream reads batches from non-seekable pipes and sockets without hanging or raising position errors while preserving streaming behavior.

Written by the indexing model from the issue text.

Assessment

Tech stack
julia
Domain
data
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
42/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.