Commit 3a0decd0 by chenyuanjie

fix

parent 8d3c2240
...@@ -52,10 +52,7 @@ class KafkaFlowAsinDetail(Templates): ...@@ -52,10 +52,7 @@ class KafkaFlowAsinDetail(Templates):
# 富集策略相关配置,用于更新 usr_mask_type 字段 # 富集策略相关配置,用于更新 usr_mask_type 字段
self.policy_name1 = "user_mask_asin_policy" self.policy_name1 = "user_mask_asin_policy"
self.policy_name2 = "user_mask_category_policy" self.policy_name2 = "user_mask_category_policy"
if self.site_name == 'us': self.pipeline_id = "user_asin_mask_enrich_pipeline"
self.pipeline_id = f"{self.site_name}_user_mask_and_profit_rate_pipeline"
else:
self.pipeline_id = ""
self.es_options = EsUtils.get_es_options(self.es_index_name, self.pipeline_id) self.es_options = EsUtils.get_es_options(self.es_index_name, self.pipeline_id)
self.db_save = 'kafka_flow_asin_detail' self.db_save = 'kafka_flow_asin_detail'
self.app_name = self.get_app_name() self.app_name = self.get_app_name()
...@@ -837,11 +834,11 @@ class KafkaFlowAsinDetail(Templates): ...@@ -837,11 +834,11 @@ class KafkaFlowAsinDetail(Templates):
EsUtils.create_index(self.es_index_name, self.client, self.es_index_body) EsUtils.create_index(self.es_index_name, self.client, self.es_index_body)
print("索引名称为:", self.es_index_name) print("索引名称为:", self.es_index_name)
# 执行富集策略 # 执行富集策略
if self.site_name == 'us': # if self.site_name == 'us':
self.client.enrich.execute_policy(name=self.policy_name1) # self.client.enrich.execute_policy(name=self.policy_name1)
self.client.enrich.execute_policy(name=self.policy_name2) # self.client.enrich.execute_policy(name=self.policy_name2)
else: # else:
pass # pass
# EsUtils.user_enrich_pipeline(self.client, self.pipeline_id, self.policy_name1, self.policy_name2) # EsUtils.user_enrich_pipeline(self.client, self.pipeline_id, self.policy_name1, self.policy_name2)
# if not EsUtils.exist_index_alias(self.es_index_alias_name, self.client): # if not EsUtils.exist_index_alias(self.es_index_alias_name, self.client):
# EsUtils.create_index_alias(self.es_index_name, self.es_index_alias_name, self.client) # EsUtils.create_index_alias(self.es_index_name, self.es_index_alias_name, self.client)
......
...@@ -51,10 +51,7 @@ class KafkaRankAsinDetail(Templates): ...@@ -51,10 +51,7 @@ class KafkaRankAsinDetail(Templates):
# 富集策略相关配置,用于更新 usr_mask_type 字段 # 富集策略相关配置,用于更新 usr_mask_type 字段
self.policy_name1 = "user_mask_asin_policy" self.policy_name1 = "user_mask_asin_policy"
self.policy_name2 = "user_mask_category_policy" self.policy_name2 = "user_mask_category_policy"
if self.site_name == 'us': self.pipeline_id = "user_asin_mask_enrich_pipeline"
self.pipeline_id = f"{self.site_name}_user_mask_and_profit_rate_pipeline"
else:
self.pipeline_id = ""
self.es_options = EsUtils.get_es_options(self.es_index_name, self.pipeline_id) self.es_options = EsUtils.get_es_options(self.es_index_name, self.pipeline_id)
self.db_save = 'kafka_rank_asin_detail' self.db_save = 'kafka_rank_asin_detail'
self.app_name = self.get_app_name() self.app_name = self.get_app_name()
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment