airbytehq / airbytehq/airbyte-python-cdk

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

未关闭
#1,129 1 条评论 0 个 reaction 已指派 0 人 在 GitHub 查看
community
主要语言
Python
星标
26
派生
53
平均合并
2 天 6 小时
30 天内合并 PR
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 摘要。