from models import ExtractorResource
class MockExtractor:
def init(self, integration_metadata: dict):
self.integration_metadata = integration_metadata
def get_resources(self):
for idx in range(5):
yield ExtractorResource(path=f'mock_s3_file_{idx}.csv')
ps: можно заматить, экстрактор реализован как генератор (yield) - в конкретном случае не так важно, тк экстрактор не отдает сами данные, а только ссылки (путь на S3 или url и тд)
Остальные классы, реализуем самостоятельно или заглядываем в репу
Каждый из классов имеет единственный аргумент = содержимое yml конфига и ничего более (но это не точно 😉), то есть все необходимые параметры должны быть описаны
в теле самого yml в нужной секции, например, у меня получилось так для экстрактора:
tasks:
extractor:
MockExtractor:
src_s3_conection_id: reddit_s3_connection_id
src_s3_bucket: raw-public
src_s3_prefix_template: reddit/{dm_date}
src_s3_partition_fmt: '%Y-%m-%d'
Осталось встроить наши классы в даг, используем TaskFlow API, пример для экстратора:
@task
def _extractor(extractor: t.Callable) -> t.List[str]:
extractor_obj = extractor(
intergation_metadata=intergation_metadata
)
return [resource.dict for resource in extractor_obj.get_resources()]
# извлекаем общую структуру тасок
tasks_meta = intergation_metadata.get("tasks", {})
# извлекаем extractor
extractor_name, extractor_params = list(tasks_meta.get("extractor").items())[0]
if not extractor_params:
extractor_params = {}
logging.info([extractor_name, extractor_params])
extractor = globals()[extractor_name]
# возвращаем список объектов
ext_resources = extractor.override(task_id=f"extractor{extractor_name}")(
extractor=extractor)
Комментарии:
- декораторная таска принимает только саму фукнцию или класс в нашем случае
- имя экстратора получаем из ямла
- все экстраторы импортируются
from extractors import *- из словаря
globals() получаем объект нужного экстрактора по имени из ямла- если параметров экстратора нет, то подставляют пустой словарь
-
resource.dict - нужно, тк XComm не знает как сериализовать нашу модель ресурса, поэтому воспользуемся атрибутом dict. Соответственно внутри таски transform_and_save обратно создадим ресурсЛогика для трансформера и saver сохраняется, добавляется только обработка ситуации,
когда трансформер сохраняет сам объекты и отдает только пути (пустой
TransformerResource.content). Как показала практика: такое нередко встречается.
Остальное обдумываем сами или заглядываем в репу.
И после загрузки в AirFlow получаем красоту в UI - repo
#walle
#framework
#automate