第2章 Dagster Resource 最佳实践:打造干净的数据管道代码
第2章 Dagster Resources 最佳实践:让数据管道的代码更干净
依赖注入和聪明的资源管理,如何拯救你的 sanity(以及你的部署)。我们经常说 Dagster 把软件工程的最佳实践带到了数据工程领域。Resource(资源)就是这样一种抽象:帮你避免重复代码,处理依赖注入,方便测试,并构建模块化、可扩展的数据平台。
什么是 Resource
Resource 是在 Dagster 中定义外部服务、工具和存储位置的标准方式。不用把每个数据库连接、API client 或存储位置硬编码进 asset,你只需把它们定义成 resource 一次,然后在需要的地方注入。
看一个简单的例子:
import dagster as dg
class DatabaseResource(dg.ConfigurableResource): connection_string: str def get_connection(self): # In real life, this would return a proper connection return f"Connected to {self.connection_string}"
@dg.asset def user_data(database: DatabaseResource): conn = database.get_connection() # Your actual logic here return "user data"
只要 Python 能访问到的东西,都可以做成 resource——数据库、API、文件系统、遗留系统统统没问题。
环境管理
管理不同环境是数据工程中比较棘手的问题之一,这主要归咎于云系统。无论是数据仓库还是分布式处理系统,要在本地、staging 和生产环境获得一致的体验都很困难。针对数据仓库,很多数据工程师的做法是为不同环境使用不同的 database 和 schema,这样无论在哪个环境,体验都保持一致。
Dagster 是这样处理的:
import dagster as dg import os
class SnowflakeResource(dg.ConfigurableResource): account: str user: str password: str database: str schema: str warehouse: str @property def connection_params(self): return { 'account': self.account, 'user': self.user, 'password': self.password, 'database': self.database, 'schema': self._get_schema(), 'warehouse': self.warehouse } def _get_schema(self): # Different schema based on environment env = os.getenv('DAGSTER_ENVIRONMENT', 'local') if env == 'prod': return self.schema elif env == 'staging': return f"{self.schema}_staging" else: return f"{self.schema}_{env}"
这样一来,asset 不需要知道自己运行在哪个环境——它们只管用 resource,剩下的交给 resource 处理:
@dg.asset def daily_metrics(snowflake: SnowflakeResource): # This automatically goes to the right place query = "SELECT * FROM metrics WHERE date = CURRENT_DATE" # Execute against whatever environment we're in return execute_query(snowflake.connection_params, query)
依赖注入:数据管道的干净架构
Resource 在依赖注入方面才真正展现威力。你的 asset 只需声明自己需要什么,而不必了解数据库连接、API key 和文件路径的细节。这种分离让业务逻辑保持干净、专注于数据转换,同时基础设施相关的问题单独处理。
不用 resource,你的代码库里可能散落着这样的东西:
# Don't do this @dg.asset def messy_asset(): # Hard-coded nightmare conn = snowflake.connector.connect( account='xy12345.us-east-1', user='data_eng_user', password='definitely_not_in_version_control', database='ANALYTICS_DB', schema='STAGING_SCHEMA', warehouse='COMPUTE_WH' ) api_client = requests.Session() api_client.headers.update({ 'Authorization': 'Bearer another_secret_token', 'User-Agent': 'our-data-pipeline/1.0' }) # And then your actual business logic gets lost in the noise data = api_client.get('https://api.example.com/data') # ... do something with data return processed_data
用了 resource,你的 asset 就干净聚焦了:
@dg.asset def clean_asset( database: SnowflakeResource, api_client: APIResource ): # Clear, focused business logic raw_data = api_client.fetch_data() processed_data = transform_data(raw_data) database.save(processed_data) return processed_data
API 封装
resource 最好用的场景之一是封装 REST API。这是来自我们 Scout 集成的真实例子(它为 Slack 频道和文档中的 ask-ai 功能提供支持):
import dagster as dg import requests from typing import Dict, List
class ScoutResource(dg.ConfigurableResource): api_key: str base_url: str = "https://api.scout.example.com" def _get_headers(self) -> Dict[str, str]: return { 'Authorization': f'Bearer {self.api_key}', 'Content-Type': 'application/json' } def write_documents(self, documents: List[Dict]) -> bool: """Upload documents to Scout for indexing""" response = requests.post( f"{self.base_url}/documents", json={"documents": documents}, headers=self._get_headers() ) return response.status_code == 200 def search_documents(self, query: str) -> List[Dict]: """Search indexed documents""" response = requests.get( f"{self.base_url}/search", params={"q": query}, headers=self._get_headers() ) return response.json().get('results', [])
使用 Scout 的 asset 现在清爽多了:
@dg.asset def index_documentation(scout: ScoutResource): docs = load_documentation_from_somewhere() success = scout.write_documents(docs) return {"documents_indexed": len(docs), "success": success}
@dg.asset def search_results(scout: ScoutResource): results = scout.search_documents("dagster best practices") return process_search_results(results)
无痛测试
测试是软件工程最佳实践中,技术功底较浅的数据从业者最容易忽略的一项。引入一些简单的功能测试,可以大幅提升数据平台的可靠性。不过测试要讲策略——没人会为了验证转换逻辑而专门启动一个 Snowflake 实例。有了 resource,你可以创建轻量的 mock,返回可预测的测试数据:
class MockDatabaseResource(ConfigurableResource): def get_users(self): return [ {"id": 1, "name": "Alice", "email": "alice@example.com"}, {"id": 2, "name": "Bob", "email": "bob@example.com"} ] def save_processed_data(self, data): # In tests, just verify the data structure assert isinstance(data, list) assert all('processed_at' in item for item in data) return True
def test_user_processing(): # Use the mock instead of real database result = user_processing_asset(MockDatabaseResource()) assert len(result) == 2 assert result[0]['processed_at'] is not None
测试从几分钟缩短到几毫秒,也不再依赖外部服务。你测的才是真正重要的东西:业务逻辑和数据转换。
我们最近在 Dagster University 上线了一门关于 Dagster 测试的课程。如果你刚接触数据工程中的测试,或者想学习 Dagster 的测试最佳实践,它是个很好的资源。
配置一目了然
Resource 的配置会展示在 Dagster UI 的 deployment 标签页下,方便你查看哪些 asset 用了哪些 resource 以及它们是如何配置的。在排查问题或给新成员介绍项目时,这种可见性非常关键。
# In your definitions file import dagster as dg
defs = dg.Definitions( assets=[user_data, daily_metrics, search_results], resources={ "database": SnowflakeResource( account="xy12345.us-east-1", user="data_eng_user", password={"env": "SNOWFLAKE_PASSWORD"}, database="ANALYTICS_DB", schema="PRODUCTION", warehouse="COMPUTE_WH" ), "scout": ScoutResource( api_key={"env": "SCOUT_API_KEY"} ) } )
什么时候该用 Resource
凡是需要和 API、数据库交互,或者项目中反复出现某种通用模式的场景,都可以用 resource。如果你在用 Fivetran、Snowflake 这类 Dagster 官方集成,通常也需要配置一个 resource 来填入 API key 和连接信息。
经验法则:如果你发现自己在多个 asset 里重复同样的初始化代码,那就是一个该抽象成 resource 的信号。
更大的意义
Resource 属于那种既解决实际问题、又能让你成为更好工程师的实践。它强制你干净地分离关注点,让代码更可测试,还让你免受硬编码配置带来的各种环境问题。
更重要的是,它能随团队一起扩展。当多人协作同一个项目时,resource 提供了与外部服务交互的标准方式。新成员不需要摸索怎么连接数据库——直接用 resource 就行。
开始上手
从小处着手。挑一个你目前硬编码的外部服务,把它改造成 resource。你会立刻在代码清晰度和可测试性上尝到甜头,然后再逐步推广。
记住:只要 Python 能访问,就能做成 resource。相信我们,花时间把这件事做对,未来的你(和你的队友)都会感谢现在的自己。
有反馈或问题?欢迎在 Slack 或 Github 上发起讨论。想和我们一起工作?看看我们的在招职位。想要更多这类内容?在 LinkedIn 上关注我们。