airbytehq / airbytehq/airbyte-python-cdk

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

オープン
#1,129 コメント 1 件 リアクション 0 件 担当者 0 名 GitHub で見る
community
主要言語
Python
スター
26
フォーク
53
平均マージ
2日 6時間
マージ済み PR(30日)
10

説明

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

コントリビューションガイド

コントリビューションガイドを開く

調査の方向性

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.

索引モデルが issue の本文から書いたものです。

評価

技術スタック
python
領域
backend, networking
issue の種類
バグ
難易度
4/5
見積もり時間
3〜5日
活発さ
活発
明瞭さ
おおむね明確
初心者へのやさしさ
48/100

新しい issue をメールで受け取る

初心者向けの GitHub issue を短くまとめたダイジェスト。