Dagster中如何将TypedDict作为Op输出处理?
问题描述
为了在不同Op中追踪Dataframe的来源,定义了TypedDict类型TaggedDF:
class TaggedDF(TypedDict): df: DataFrame ressource_name: str
并编写了如下Op:
@op(retry_policy=OPS_RETRY_POLICY) def read_fec(context, file_name: str) -> TaggedDF: """ read fec files from source container """ context.log.info("reading %s", file_name) reg = match(r".*\/(.*)FEC(.*)\.txt", file_name) blob_client_instance = azure_client.get_blob_client( Vars.LANDING_ZONE.value, file_name ) blob_data = blob_client_instance.download_blob() df = pd.read_csv(StringIO(blob_data.content_as_text(encoding="UTF-8")), sep="\t") df["siren"] = reg.group(1) df["period"] = reg.group(2) return {"df": df, "ressource_name": file_name}
但运行时触发Dagster类型检查错误:
dagster._core.errors.DagsterTypeCheckError: Error occurred while type-checking output "result" of op "read_fec", with Python type <class 'dict'> and Dagster type TaggedDF File "/home/zar3bski/.cache/pypoetry/virtualenvs/poc-dagster-30bqhDW5-py3.10/lib/python3.10/site-packages/dagster/_core/execution/plan/execute_plan.py", line 269, in dagster_event_sequence_for_step for step_event in check.generator(step_events): File "/home/zar3bski/.cache/pypoetry/virtualenvs/poc-dagster-30bqhDW5-py3.10/lib/python3.10/site-packages/dagster/_core/execution/plan/execute_step.py", line 386, in core_dagster_event_sequence_for_step for evt in _type_check_and_store_output(step_context, user_event): File "/home/zar3bski/.cache/pypoetry/virtualenvs/poc-dagster-30bqhDW5-py3.10/lib/python3.10/site-packages/dagster/_core/execution/plan/execute_step.py", line 436, in _type_check_and_store_output for output_event in _type_check_output(step_context, step_output_handle, output, version): File "/home/zar3bski/.cache/pypoetry/virtualenvs/poc-dagster-30bqhDW5-py3.10/lib/python3.10/site-packages/dagster/_core/execution/plan/execute_step.py", line 282, in _type_check_output with user_code_error_boundary( File "/usr/lib/python3.10/contextlib.py", line 153, in __exit__ self.gen.throw(typ, value, traceback) File "/home/zar3bski/.cache/pypoetry/virtualenvs/poc-dagster-30bqhDW5-py3.10/lib/python3.10/site-packages/dagster/_core/errors.py", line 213, in user_code_error_boundary raise error_cls( The above exception was caused by the following exception: TypeError: TypedDict does not support instance and class checks File "/home/zar3bski/.cache/pypoetry/virtualenvs/poc-dagster-30bqhDW5-py3.10/lib/python3.10/site-packages/dagster/_core/errors.py", line 206, in user_code_error_boundary yield File "/home/zar3bski/.cache/pypoetry/virtualenvs/poc-dagster-30bqhDW5-py3.10/lib/python3.10/site-packages/dagster/_core/execution/plan/execute_step.py", line 287, in _type_check_output type_check = do_type_check(type_check_context, dagster_type, output.value) File "/home/zar3bski/.cache/pypoetry/virtualenvs/poc-dagster-30bqhDW5-py3.10/lib/python3.10/site-packages/dagster/_core/execution/plan/execute_step.py", line 188, in do_type_check type_check = dagster_type.type_check(context, value) File "/home/zar3bski/.cache/pypoetry/virtualenvs/poc-dagster-30bqhDW5-py3.10/lib/python3.10/site-packages/dagster/_core/types/dagster_type.py", line 171, in type_check retval = self._type_check_fn(context, value) File "/home/zar3bski/.cache/pypoetry/virtualenvs/poc-dagster-30bqhDW5-py3.10/lib/python3.10/site-packages/dagster/_core/types/dagster_type.py", line 507, in type_check if not isinstance(value, expected_python_type): File "/usr/lib/python3.10/typing.py", line 2385, in __subclasscheck__ raise TypeError('TypedDict does not support instance and class checks')
解决方案
Dagster无法直接用Python的TypedDict作为Op输出类型,核心原因是TypedDict不支持isinstance检查,而Dagster的类型验证依赖该操作。可通过以下两种方式解决:
方法一:自定义Dagster类型
创建自定义Dagster类型,手动实现类型检查逻辑:
from dagster import DagsterType, TypeCheckContext from pandas import DataFrame def _tagged_df_type_check(context: TypeCheckContext, value: object) -> bool: # 先验证是否为字典 if not isinstance(value, dict): return False # 验证是否包含必填键 required_keys = {"df", "ressource_name"} if not required_keys.issubset(value.keys()): return False # 验证字段类型 return isinstance(value["df"], DataFrame) and isinstance(value["ressource_name"], str) # 注册自定义Dagster类型 TaggedDF = DagsterType( name="TaggedDF", type_check_fn=_tagged_df_type_check, description="包含DataFrame及其来源资源名称的字典" )
之后在Op中直接使用该自定义类型作为返回类型即可,原Op的返回字典逻辑无需修改。
方法二:用dataclass替代TypedDict
改用Python的dataclass定义结构,Dagster可自动支持其类型检查:
from dataclasses import dataclass from pandas import DataFrame @dataclass class TaggedDF: df: DataFrame ressource_name: str
修改Op的返回逻辑,返回TaggedDF实例而非字典:
@op(retry_policy=OPS_RETRY_POLICY) def read_fec(context, file_name: str) -> TaggedDF: # 原业务逻辑保持不变 ... return TaggedDF(df=df, ressource_name=file_name)
这种写法更符合面向对象风格,同时无需额外编写类型检查逻辑,Dagster会自动处理类型验证。
内容的提问来源于stack exchange,提问作者zar3bski
相关产品推荐
相关产品推荐

