
dlt 的 rest_api_source 本身支持声明式并行化,无需手动封装 @dlt.resource(parallelized=True);只需在资源定义中设置 parallelize=True(注意拼写为 parallelize,非 parallelized),即可让每个资源实例独立并发执行。
在 dlt 中正确启用 rest api 管道的并行化处理,关键在于理解 `rest_api_source` 的原生并行机制——它并非通过装饰器函数实现,而是直接在资源级配置中启用 `parallelize=true` 参数。错误地尝试用 `@dlt.resource(parallelized=true)` 包裹动态生成的资源对象,不仅违背 dlt rest api 源的设计范式,还会触发 `resourcenamemissing` 等类型校验异常。
dlt 官方文档明确指出:rest_api_source 是一个声明式、可并行的源(declarative & parallelizable source)。其底层已内置对多资源并发拉取的支持,只需在每个资源字典中显式添加 "parallelize": True 字段,dlt 运行时便会自动为该资源启动独立线程/进程(取决于执行器配置),而无需重构为函数式资源。
✅ 正确做法(推荐):
在原始 resources 列表推导式中,为每个资源添加 "parallelize": True:
from dlt.sources.rest_api import rest_api_source
source = rest_api_source({
"client": {
"base_url": "https://www.filmweb.no/",
# 可选:配置并发数(默认为 CPU 核心数)
"auth": None,
},
"resources": [
{
"name": f"movie_{MOVIE_ID}",
"table_name": "movies",
"parallelize": True, # ← 关键:启用该资源的并行执行
"endpoint": {
"path": "_next/data/{build_id}/film/{movie_id}.json",
"params": {
"movie_id": MOVIE_ID,
"build_id": BUILD_ID,
"edi": MOVIE_ID,
},
"data_selector": "pageProps.cmsDocument",
},
"write_disposition": "replace",
}
for MOVIE_ID in MOVIE_IDS
],
})
pipeline = dlt.pipeline(pipeline_name="movies", destination="filesystem")
load_info = pipeline.run(source)⚠️ 注意事项:
parallelize是资源级开关(布尔值),不是@dlt.resource的参数;parallelized=True是旧版或误传写法,当前稳定版(v1.0+)仅识别parallelize。所有并行资源共享同一
client配置(如base_url,auth,timeout),但各自拥有独立的请求上下文与重试策略。-
若需精细控制并发度(如限制最大并发请求数),可在
pipeline.run()时传入workers=8参数(需配合threading或processes执行器):load_info = pipeline.run(source, workers=4) # 最多 4 个并发请求
-
建议使用类型提示
RESTAPIConfig(来自dlt.sources.rest_api)提升 IDE 自动补全与类型安全:from dlt.sources.rest_api import RESTAPIConfig config: RESTAPIConfig = { ... } # IDE 将提示可用字段
? 总结:dlt 的 REST API 并行化是“开箱即用”的声明式能力,核心在于资源定义中的 parallelize=True。避免将 rest_api_source 嵌套于自定义函数或资源装饰器中——这会破坏其内部资源注册与调度逻辑,导致元数据缺失(如 ResourceNameMissing)。保持配置简洁、语义清晰,才是高效利用 dlt 并行能力的最佳实践。


















