Airbyte自定义连接器:同一类中如何复用函数返回值?
问题解答
首先明确核心限制:Airbyte中get_json_schema()的调用时机早于parse_response()——系统会在启动数据同步之前先调用get_json_schema()获取流的结构定义,而parse_response()是在拉取到API数据后才执行的方法。所以直接复用parse_response()里的response_json变量是不可行的,因为前者执行时后者还没产生数据。
针对你的场景,给两个可行方案:
方案1:在get_json_schema()中主动预取Schema数据
直接在get_json_schema()里调用API获取样本数据(比如请求单条记录),用这个数据生成Schema,和parse_response()复用相同的API请求逻辑即可。示例伪代码:
class SurveyctoStream(SurveyStream): def _fetch_sample_data(self): # 复用和parse_response一致的API请求逻辑,获取一条样本数据 response = self._send_request() return response.json()[0] if response.json() else {} def get_json_schema(self): sample_data = self._fetch_sample_data() # 从样本数据生成JSON Schema(可以用jsonschema库的infer_schema方法) return infer_schema(sample_data) def parse_response(self, response, **kwargs): response_json = response.json() # 处理数据逻辑 return response_json
注意:要给API请求加缓存(比如用lru_cache装饰器),避免每次调用get_json_schema()都发起API请求,降低性能消耗。
方案2:新增工具类解耦职责
如果你的Schema生成逻辑比较复杂,或者需要在多个地方复用,可以新增一个SurveySchemaGenerator类,专门负责从API拉取数据并生成Schema:
class SurveySchemaGenerator: def __init__(self, config, form_id): self.config = config self.form_id = form_id def fetch_sample_data(self): # 封装API请求逻辑 response = requests.get(...) return response.json()[0] if response.json() else {} def generate_schema(self): sample_data = self.fetch_sample_data() return infer_schema(sample_data) class SurveyctoStream(SurveyStream): def __init__(self, config, form_id): super().__init__(config) self.form_id = form_id self.schema_generator = SurveySchemaGenerator(config, form_id) def get_json_schema(self): return self.schema_generator.generate_schema() def parse_response(self, response, **kwargs): response_json = response.json() # 处理数据逻辑 return response_json
这种方式把Schema生成和数据拉取的职责分开,代码更清晰,也方便后续扩展。
额外提醒:因为你是针对每个form_id创建独立流,要确保每个流实例初始化时,都对应正确的form_id来生成专属Schema,避免不同表单的Schema混淆。
内容的提问来源于stack exchange,提问作者Siddhant Singh
相关产品推荐
相关产品推荐

