airbytehq / airbytehq/airbyte-python-cdk

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

Aberta
#1,129 1 comentário 0 reações 0 responsáveis Ver no GitHub
community
Linguagem predominante
Python
Estrelas
26
Forks
53
Merge médio
2d 6h
PRs com merge (30d)
10

Descrição

## 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.

Guia de contribuição

Abrir o guia de contribuição

Direção de pesquisa

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.

Escrita pelo modelo de indexação a partir do texto da issue.

Avaliação

Stack de tecnologia
python
Domínio
backend, networking
Tipo de issue
Bug
Dificuldade
4/5
Tempo estimado
3-5 dias
Status de atividade
Ativa
Clareza
Razoavelmente clara
Facilidade para iniciantes
48/100

Receba novas issues na sua caixa de entrada

Um resumo curto de issues do GitHub para quem está começando.