Skip to content

flight

Arrow Flight utilities.

ChainedFlightReader

ChainedFlightReader(endpoints, initial_url, metadata=None, middleware=None)

Bases: StreamReader

Sequential reader for time-ordered endpoint slices.

Groups endpoints by time slice (based on the start time in their tickets), orders slices chronologically, and reads each slice to completion before starting the next. Concurrent endpoints within the same slice are handled by a MultiFlightReader.

Source code in arrakis/flight.py
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
def __init__(
    self,
    endpoints: list[flight.FlightEndpoint],
    initial_url: str | flight.FlightClient,
    metadata: dict[str, Channel] | None = None,
    middleware: list[flight.ClientMiddlewareFactory] | None = None,
):
    self.initial_url = initial_url
    self.metadata = metadata
    self._middleware = middleware or []

    # group endpoints by start time
    slices: dict[int | None, list[flight.FlightEndpoint]] = {}
    for endpoint in endpoints:
        start = _endpoint_start_time(endpoint)
        slices.setdefault(start, []).append(endpoint)

    # order slices chronologically (None sorts first as live/unbounded)
    self._slices = sorted(
        slices.items(),
        key=lambda item: item[0] if item[0] is not None else float("inf"),
    )

    # single slice: delegate directly to one MultiFlightReader
    # and avoid the chaining/busy-wait overhead entirely
    if len(self._slices) == 1:
        _, endpoints = self._slices[0]
        self._delegate = self._make_reader(
            endpoints, self.initial_url, self.metadata
        )
    else:
        self._delegate = None

    self._current_reader: MultiFlightReader | None = None
    self._slice_index = 0

MultiFlightReader

MultiFlightReader(endpoints, initial_url, metadata=None, middleware=None)

Bases: StreamReader

Multi-threaded Arrow Flight endpoint stream iterator context manager

Given a list of endpoints, connect to all of them in parallel and stream data from them all interleaved.

Source code in arrakis/flight.py
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
def __init__(
    self,
    endpoints: list[flight.FlightEndpoint],
    initial_url: str | flight.FlightClient,
    metadata: dict[str, Channel] | None = None,
    middleware: list[flight.ClientMiddlewareFactory] | None = None,
):
    """initialize with list of endpoints and an reusable flight client"""
    self.endpoints = endpoints
    self.initial_url = initial_url
    self.metadata = metadata
    self._middleware = middleware or []

    self.queue: queue.SimpleQueue = queue.SimpleQueue()
    # FIXME: this quit event could be replaced with the
    # Queue.shutdown method after python 3.13
    self.quit_event: threading.Event = threading.Event()
    self.done_events = [threading.Event() for _ in self.endpoints]
    self.executor = concurrent.futures.ThreadPoolExecutor(
        max_workers=len(self.endpoints),
    )
    self.futures: list[concurrent.futures.Future] = []
    # in-flight DoGet readers, so close() can cancel a read
    # blocked on a quiet stream instead of waiting on it forever
    self._readers: list[flight.FlightStreamReader] = []
    self._readers_lock = threading.Lock()
    # per-endpoint liveness watchdog (see _raise_stalled_endpoints):
    # wall-clock silence bound and last delivery time per stream id
    self._stall_bounds = self._stall_watchdog_bounds()
    self._last_delivery: dict[str, float] = {}
    self._future_by_id: dict[str, concurrent.futures.Future] = {}

close

close()

Close all streams, cancelling in-flight reads.

A DoGet blocked on a quiet stream holds its worker thread until data arrives; without the cancel, close() waited on that forever (and interrupting it left the interpreter's atexit hook waiting on the same thread). Cancelling unblocks the read immediately and signals the server, whose stream generator observes the cancellation at its next heartbeat.

Source code in arrakis/flight.py
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
def close(self):
    """Close all streams, cancelling in-flight reads.

    A DoGet blocked on a quiet stream holds its worker thread
    until data arrives; without the cancel, close() waited on
    that forever (and interrupting it left the interpreter's
    atexit hook waiting on the same thread).  Cancelling unblocks
    the read immediately and signals the server, whose stream
    generator observes the cancellation at its next heartbeat.
    """
    self.quit_event.set()
    with self._readers_lock:
        readers = list(self._readers)
    for reader in readers:
        with contextlib.suppress(Exception):
            reader.cancel()
    for f in self.futures:
        try:
            f.result(timeout=CLOSE_TIMEOUT)
        except flight.FlightCancelledError:
            # our own cancellation winding down
            pass
        except concurrent.futures.TimeoutError:
            logger.warning(
                "endpoint stream did not stop within %.0fs of "
                "cancellation; abandoning its worker thread",
                CLOSE_TIMEOUT,
            )
        except flight.FlightError as e:
            # NOTE: this strips the original message of everything
            # besides the original error message raised by the server
            msg = e.args[0].partition(" Detail:")[0]
            raise type(e)(msg, e.extra_info) from None

    self.executor.shutdown(wait=False, cancel_futures=True)
    self.futures = []

done

done()

returns True when all threads are done and the queue is empty

A failed endpoint must surface even here: consumers check done() before read(), so with the queue drained a stream whose endpoint(s) died with an error would otherwise be reported as a clean end of stream and the error never raised.

Source code in arrakis/flight.py
528
529
530
531
532
533
534
535
536
537
538
539
def done(self) -> bool:
    """returns True when all threads are done and the queue is empty

    A failed endpoint must surface even here: consumers check
    done() before read(), so with the queue drained a stream whose
    endpoint(s) died with an error would otherwise be reported as
    a clean end of stream and the error never raised.
    """
    if all(e.is_set() for e in self.done_events) and self.queue.empty():
        self._raise_failed_endpoints()
        return True
    return False

read

read(*, convert_blocks=False)

read all available data from all streams

Yields:

Type Description
tuple of (stream_id, data)
Source code in arrakis/flight.py
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
def read(
    self,
    *,
    convert_blocks: bool = False,
) -> Generator[
    tuple[str, Any],
    None,
    None,
]:
    """read all available data from all streams

    Yields
    ------
    tuple of (stream_id, data)

    """
    if convert_blocks:
        assert self.metadata, (
            "Must be initialized with channel metadata to convert to blocks."
        )
    while True:
        try:
            stream_id, batch = self.queue.get(block=False)
            if convert_blocks:
                assert self.metadata  # needed for type checking consistency
                data = SeriesBlock.from_column_batch(
                    batch,
                    self.metadata,
                )
            else:
                data = batch
            yield stream_id, data
        except queue.Empty:
            break
    self._raise_failed_endpoints()
    self._raise_stalled_endpoints()

ParsedRequest dataclass

ParsedRequest(request, args, version)

A parsed and validated request envelope.

RequestValidator

RequestValidator()

A validator for JSON-encoded requests.

Source code in arrakis/flight.py
 99
100
101
102
103
104
105
def __init__(self) -> None:
    self._schemas: dict[RequestType, dict[str, Any]] = {}
    self._validators: dict[RequestType, jsonschema.Draft7Validator] = {}

    # load generic descriptor schema
    schema = arrakis_schema.load_schema("descriptor.json")
    self._generic_validator = jsonschema.Draft7Validator(schema)

known_args

known_args(request)

Return the argument names the loaded schema defines for a request.

Servers use this to distinguish arguments introduced by newer schema versions from typos, and reject both with an actionable error instead of an incidental dispatch failure.

Source code in arrakis/flight.py
115
116
117
118
119
120
121
122
def known_args(self, request: RequestType) -> frozenset[str]:
    """Return the argument names the loaded schema defines for a request.

    Servers use this to distinguish arguments introduced by newer
    schema versions from typos, and reject both with an actionable
    error instead of an incidental dispatch failure.
    """
    return frozenset(self._schema(request)["properties"]["args"]["properties"])

validate

validate(payload)

Validate a JSON-encoded request.

Parameters:

Name Type Description Default
payload Request

A dictionary with a 'request' and an 'args' key encoding the given Flight request.

required

Raises:

Type Description
ValidationError

If the request does not match the expected schema.

Source code in arrakis/flight.py
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
def validate(self, payload: Request) -> None:
    """Validate a JSON-encoded request.

    Parameters
    ----------
    payload : Request
        A dictionary with a 'request' and an 'args' key encoding
        the given Flight request.

    Raises
    ------
    ValidationError
        If the request does not match the expected schema.

    """
    self._generic_validator.validate(payload)
    request = RequestType[payload["request"]]

    if request not in self._validators:
        self._validators[request] = jsonschema.Draft7Validator(
            self._schema(request)
        )

    self._validators[request].validate(payload)

connect

connect(location, **kwargs)

Connect to a Flight server with arrakis default channel options.

A drop-in wrapper for :func:pyarrow.flight.connect that applies KEEPALIVE_OPTIONS unless the caller passes its own generic_options.

Source code in arrakis/flight.py
76
77
78
79
80
81
82
83
84
def connect(location, **kwargs) -> flight.FlightClient:
    """Connect to a Flight server with arrakis default channel options.

    A drop-in wrapper for :func:`pyarrow.flight.connect` that applies
    ``KEEPALIVE_OPTIONS`` unless the caller passes its own
    ``generic_options``.
    """
    kwargs.setdefault("generic_options", list(KEEPALIVE_OPTIONS))
    return flight.connect(location, **kwargs)

create_command

create_command(request_type, *, validator, **kwargs)

Create a Flight command containing a JSON-encoded request.

Parameters:

Name Type Description Default
request_type RequestType

The type of request.

required
validator RequestValidator

A validator to validate that the command matches the expected schema.

required
**kwargs dict

Extra arguments corresponding to the specific request.

{}

Returns:

Type Description
bytes

The JSON-encoded request.

Raises:

Type Description
ValidationError

If the request does not match the expected schema.

Source code in arrakis/flight.py
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
def create_command(
    request_type: RequestType, *, validator: RequestValidator, **kwargs
) -> bytes:
    """Create a Flight command containing a JSON-encoded request.

    Parameters
    ----------
    request_type : RequestType
        The type of request.
    validator : RequestValidator
        A validator to validate that the command matches the expected schema.
    **kwargs : dict, optional
        Extra arguments corresponding to the specific request.

    Returns
    -------
    bytes
        The JSON-encoded request.

    Raises
    ------
    ValidationError
        If the request does not match the expected schema.

    """
    cmd: Request = {
        "request": request_type.name,
        "args": kwargs,
        # the request schema version, so servers can report actionable
        # errors for requests newer than they support
        "version": arrakis_schema.SCHEMA_VERSION,
    }
    validator.validate(cmd)
    return json.dumps(cmd).encode("utf-8")

create_descriptor

create_descriptor(request_type, *, validator, **kwargs)

Create a Flight descriptor given a request.

Parameters:

Name Type Description Default
request_type RequestType

The type of request.

required
validator RequestValidator

A validator to validate that the command matches the expected schema.

required
**kwargs dict

Extra arguments corresponding to the specific request.

{}

Returns:

Type Description
FlightDescriptor

A Flight Descriptor containing the request.

Raises:

Type Description
ValidationError

If the request does not match the expected schema.

Source code in arrakis/flight.py
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
def create_descriptor(
    request_type: RequestType, *, validator: RequestValidator, **kwargs
) -> flight.FlightDescriptor:
    """Create a Flight descriptor given a request.

    Parameters
    ----------
    request_type : RequestType
        The type of request.
    validator : RequestValidator
        A validator to validate that the command matches the expected schema.
    **kwargs : dict, optional
        Extra arguments corresponding to the specific request.

    Returns
    -------
    flight.FlightDescriptor
        A Flight Descriptor containing the request.

    Raises
    ------
    ValidationError
        If the request does not match the expected schema.

    """
    cmd = create_command(request_type, validator=validator, **kwargs)
    return flight.FlightDescriptor.for_command(cmd)

parse_command

parse_command(cmd, *, validator)

Parse a Flight command into a request.

Parameters:

Name Type Description Default
cmd bytes

The JSON-encoded request.

required
validator RequestValidator

A validator to validate that the command matches the expected schema.

required

Returns:

Name Type Description
request_type RequestType

The type of request.

kwargs dict

Arguments corresponding to the specific request.

Raises:

Type Description
JSONDecodeError

If the command does not decode to valid JSON.

ValidationError

If the request does not match the expected schema.

Source code in arrakis/flight.py
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
def parse_command(
    cmd: bytes, *, validator: RequestValidator
) -> tuple[RequestType, dict]:
    """Parse a Flight command into a request.

    Parameters
    ----------
    cmd : bytes
        The JSON-encoded request.
    validator : RequestValidator
        A validator to validate that the command matches the expected schema.

    Returns
    -------
    request_type : RequestType
        The type of request.
    kwargs : dict
        Arguments corresponding to the specific request.

    Raises
    ------
    JSONDecodeError
        If the command does not decode to valid JSON.
    ValidationError
        If the request does not match the expected schema.

    """
    parsed = parse_request(cmd, validator=validator)
    return parsed.request, parsed.args

parse_request

parse_request(cmd, *, validator)

Parse a Flight command into a request envelope.

Parameters:

Name Type Description Default
cmd bytes

The JSON-encoded request.

required
validator RequestValidator

A validator to validate that the command matches the expected schema.

required

Returns:

Type Description
ParsedRequest

The parsed request envelope.

Raises:

Type Description
JSONDecodeError

If the command does not decode to valid JSON.

ValidationError

If the request does not match the expected schema.

Source code in arrakis/flight.py
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
def parse_request(cmd: bytes, *, validator: RequestValidator) -> ParsedRequest:
    """Parse a Flight command into a request envelope.

    Parameters
    ----------
    cmd : bytes
        The JSON-encoded request.
    validator : RequestValidator
        A validator to validate that the command matches the expected schema.

    Returns
    -------
    ParsedRequest
        The parsed request envelope.

    Raises
    ------
    JSONDecodeError
        If the command does not decode to valid JSON.
    ValidationError
        If the request does not match the expected schema.

    """
    try:
        parsed = json.loads(cmd.decode("utf-8"))
    except json.JSONDecodeError as e:
        msg = "Command does not decode to valid JSON"
        raise json.JSONDecodeError(msg, e.doc, e.pos) from e
    else:
        validator.validate(parsed)
        return ParsedRequest(
            request=RequestType[parsed["request"]],
            args=parsed["args"],
            version=parsed.get("version"),
        )