from enrichment.enrich_age import EnrichAge from enrichment.enrich_feature_eng import EnrichFeatureEngineering from enrichment.enrich_geo_pelias import EnrichGeographicsPelias from enrichment.enrich_ml import EnrichMachineLearningClusters from enrichment.enrich_name import EnrichName from enrichment.enrich_rfm import EnrichRFM from enrichment.enrich_superfan import EnrichSuperfan from enrichment.enrich_unique_gender import EnrichUniqueGender ENRICHMENTS = [ EnrichAge, EnrichName, EnrichUniqueGender, EnrichRFM, EnrichFeatureEngineering, EnrichGeographicsPelias, EnrichSuperfan, EnrichMachineLearningClusters, ] def get_enrichments_dependencies(endpoint, req, event): source_attribute_ids = {enrichment.__name__: enrichment.source_attribute_ids for enrichment in ENRICHMENTS} dependencies = [] for i, enr1 in enumerate(ENRICHMENTS): for enr2 in ENRICHMENTS[i + 1 :]: if len(set(enr1.result_attribute_ids).intersection(set(enr2.source_attribute_ids))) > 0: dependencies.append([enr1.__name__, enr2.__name__]) elif len(set(enr2.result_attribute_ids).intersection(set(enr1.source_attribute_ids))) > 0: dependencies.append([enr2.__name__, enr1.__name__]) return dict(source_attribute_ids=source_attribute_ids, dependencies=dependencies) airflow_endpoints = {"airflow/enrichment-dependencies": {"query": get_enrichments_dependencies}}