airbytehq / airbytehq/airbyte-python-cdk

DeclarativePartitionFactory shares one retriever instance across worker threads despite per-thread docstring

Aperta
#1,129 1 commento 0 reazioni 0 assegnatari Vedi su GitHub
community
Lingua principale
Python
Stelle
26
Fork
53
Merge medio
2g 6h
PR unite (30g)
10

Descrizione

## Summary

`DeclarativePartitionFactory` claims to create "a retriever per thread" in its docstring, but it holds and reuses a single retriever instance for every partition. With the default `concurrency_level` of a declarative source, all worker threads share that one retriever - and for HTTP-based retrievers that means one shared `requests.Session` and unlocked mutable state used concurrently from multiple threads.

## Code refs (v7.25.1)

- `airbyte_cdk/sources/declarative/stream_slicers/declarative_partition_generator.py:44-47` - docstring: the factory exists "in order to prevent the stream instance from being shared between threads" by creating "a retriever per thread".
- `airbyte_cdk/sources/declarative/stream_slicers/declarative_partition_generator.py:50` - the factory stores ONE retriever instance and passes the same object into every `DeclarativePartition` (`:91-93` calls `self._retriever.read_records(...)` per partition, potentially from N worker threads at once).
- For custom retrievers built around `airbyte_cdk/sources/streams/http/http_client.py`: one shared `requests.Session` (`:126-180`) and the unlocked `_request_attempt_count` dict (`:137`, mutated at `:329-334`) are then shared across threads.

## Impact

- Any `CustomRetriever` following the documented contract gets concurrent `read_records` calls on one instance without any warning - easy to write thread-unsafe connector code that works in unit tests (single-threaded) and misbehaves in production (default concurrency 5+).
- Shared-session symptoms are environment-dependent (connection-pool contention, retry-counter races), which makes them hard to attribute.

## Expected

Either the implementation should match the docstring (instantiate a retriever per partition/thread, e.g. deep-copy or factory callback), or the docstring and the `CustomRetriever` documentation should state explicitly that one instance is shared across worker threads and implementations must be thread-safe.

## Precedent

Found while reviewing https://github.com/airbytehq/airbyte/pull/80306 (source-freshdesk `ticket_activities` custom retriever): 5 worker threads issue concurrent `read_records` through one `HttpClient`/`requests.Session` instance.

Guida per i contributori

Apri la guida per i contributori

Direzione di ricerca

Read airbyte_cdk/sources/declarative/stream_slicers/declarative_partition_generator.py, especially the factory and DeclarativePartition calls, then inspect http_client.py and the CustomRetriever contract. Use the documented per-thread behavior and the shared-session and retry-counter references to determine the intended ownership model. Done means the implementation and documentation consistently describe and enforce safe retriever use across worker threads.

Scritto dal modello di indicizzazione a partire dal testo della issue.

Valutazione

Stack tecnologico
python
Ambito
backend, networking
Tipo di issue
Bug
Difficoltà
4/5
Tempo stimato
3-5 giorni
Stato di attività
Attiva
Chiarezza
Abbastanza chiara
Idoneità per principianti
48/100

Ricevi le nuove issue nella tua casella

Un breve riepilogo di issue GitHub adatte ai principianti.