Skip to main content

JSONL

JSONL (JSON Lines) file and directory types for Flyte, backed by orjson for fast serialization and optional zstd compression.

pip install flyteplugins-jsonl

# For Arrow RecordBatch support
pip install 'flyteplugins-jsonl[arrow]'

Types

JsonlFile

A single JSONL file. Inherits from flyte.io.File so it works with remote storage, upload/download and the Flyte type engine out of the box.

from flyteplugins.jsonl import JsonlFile

# Async read
@env.task
async def process(f: JsonlFile):
    async for record in f.iter_records():
        print(record)

# Async write
@env.task
async def create() -> JsonlFile:
    f = JsonlFile.new_remote("data.jsonl")
    async with f.writer() as w:
        await w.write({"key": "value"})
    return f

# Sync write
@env.task
async def create_sync() -> JsonlFile:
    f = JsonlFile.new_remote("data.jsonl")
    with f.writer_sync() as w:
        w.write({"key": "value"})
    return f

JsonlDir

A directory of sharded JSONL files (part-00000.jsonl, part-00001.jsonl, etc.). Inherits from flyte.io.Dir. Supports automatic shard rotation on write and transparent cross-shard iteration on read.

from flyteplugins.jsonl import JsonlDir

# Write with automatic sharding
@env.task
async def create() -> JsonlDir:
    d = JsonlDir.new_remote("output_shards")
    async with d.writer(max_records_per_shard=10_000) as w:
        for i in range(50_000):
            await w.write({"id": i})
    return d

# Read across all shards
@env.task
async def process(d: JsonlDir):
    async for record in d.iter_records():
        print(record)

Features

Compression

Both types support zstd compression transparently via file extension. Use .jsonl.zst to enable:

# Single file
f = JsonlFile.new_remote("data.jsonl.zst")

# Sharded directory
async with d.writer(shard_extension=".jsonl.zst") as w:
    await w.write({"compressed": True})

Prefetch (JsonlDir)

When iterating over a sharded directory, the next shard is prefetched in the background to overlap network I/O with processing. This is enabled by default and can be tuned or disabled:

async for record in d.iter_records(prefetch=True, queue_size=8192):
    process(record)

queue_size is the memory safety bound on the read-ahead buffer.

Batch iteration

Both types support batched iteration for bulk processing:

# List-of-dicts batches
async for batch in d.iter_batches(batch_size=1000):
    process_batch(batch)  # list[dict]

# Arrow RecordBatches (requires pyarrow)
async for batch in d.iter_arrow_batches(batch_size=65536):
    table = pa.Table.from_batches([batch])

Sync variants are available: iter_batches_sync(), iter_arrow_batches_sync().

Error handling

All read methods accept an on_error parameter:

  • "raise" (default) -- propagate parse errors immediately
  • "skip" -- log a warning and skip corrupt lines
  • A callable (line_number: int, raw_line: bytes, exception: Exception) -> None for custom handling
async for record in f.iter_records(on_error="skip"):
    print(record)

Shard rotation

The directory writer rotates shards based on record count, byte size or both:

async with d.writer(
    max_records_per_shard=10_000,       # rotate after 10k records
    max_bytes_per_shard=256 << 20,      # or after 256 MB (default)
) as w:
    ...

Append

Opening a writer on a directory that already contains shards is safe -- the writer scans for existing part-NNNNN files and starts from the next index.

Sync vs Async

Every read/write method has both async and sync variants:

Async Sync
iter_records() iter_records_sync()
iter_batches() iter_batches_sync()
iter_arrow_batches() iter_arrow_batches_sync()
writer() writer_sync()

Examples

See examples/ for runnable scripts:

  • jsonl_file.py -- single-file read/write with compression and error handling
  • jsonl_dir.py -- sharded directory read/write, append, and compression
  • jsonl_arrow.py -- Arrow RecordBatch iteration for analytics workloads

Release files for flyteplugins-jsonl 2.10.1

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Built distribution (wheel)

Table of built distributions (wheels) for flyteplugins-jsonl 2.10.1
File Interpreter ABI Platform
flyteplugins_jsonl-2.10.1-py3-none-any.whl Python 3 none any Details

Release files / flyteplugins_jsonl-2.10.1-py3-none-any.whl

Download URL flyteplugins_jsonl-2.10.1-py3-none-any.whl
Size 12.1 kB
Tags Python 3
SHA-256 checksum
How to use checksums
3ae4d0264b5c12be33e9bbe4fe05f18a2a3cc478b9c31f68ce0d291497cd4084
BLAKE2b-256 checksum
How to use checksums
4638fdda0fd6f03283abf4eaaa86f224d3a33d71515630b97a1792594955946d
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.13.15

Release history Release notifications | RSS feed

This release

2.10.1 This release

1 release file

2.10.0

1 release file

2.9.0

1 release file

2.8.1

1 release file

2.8.0

1 release file

2.7.2

1 release file

2.7.1

1 release file

2.7.0

1 release file

2.6.13

1 release file

2.6.12

1 release file

2.6.11

1 release file

2.6.10

1 release file

2.6.9

1 release file

2.6.8

1 release file

2.6.7

1 release file

2.6.6

1 release file

2.6.5

1 release file

2.6.4

1 release file

2.6.3

1 release file

2.6.2

1 release file

2.6.1

1 release file

2.6.0

1 release file

2.5.20

1 release file

2.5.18

1 release file

2.5.17

1 release file

2.5.14

1 release file

2.5.13

1 release file

2.5.12

1 release file

2.5.11

1 release file

2.5.10

1 release file

2.5.9

1 release file

2.5.8

1 release file

2.5.7

1 release file

2.5.6

1 release file

2.5.5

1 release file

2.5.4

1 release file

2.5.3

1 release file

2.5.2

1 release file

2.5.1

1 release file

2.5.0

1 release file

2.4.4

1 release file

2.4.3

1 release file

2.4.2

1 release file

2.4.1

1 release file

2.4.0

1 release file

2.3.9

1 release file

2.3.8

1 release file

2.3.7

1 release file

2.3.6

1 release file

2.3.5

1 release file

2.3.4

1 release file

2.3.3

1 release file

2.3.2

1 release file

2.3.1

1 release file

2.3.0

1 release file

2.2.4

1 release file

2.2.3

1 release file

2.2.2

1 release file

2.2.1

1 release file

2.2.0

1 release file

2.1.9

1 release file

2.1.8

1 release file

2.1.7

1 release file

2.1.6

1 release file

2.1.5

1 release file

2.1.4

1 release file

2.1.3

1 release file

2.1.2

1 release file

2.1.1

1 release file

2.1.0

1 release file

2.0.12

1 release file

2.0.11

1 release file

2.0.10

1 release file

2.0.9

1 release file

2.0.8

1 release file

2.0.7

1 release file

2.0.6

1 release file

2.0.4

1 release file

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page