Commit 1165e91b by chenyuanjie

fix

parent 599cd475
......@@ -29,7 +29,7 @@ NEED_FILTER_CATEGORIES = (
class KafkaFlowAsinDetail(Templates):
def __init__(self, site_name='us', date_type="month", date_info='2026-03', consumer_type='history', test_flag='test', batch_size=1200000):
def __init__(self, site_name='us', date_type="month", date_info='2026-03', consumer_type='history', test_flag='test', batch_size=240000):
super().__init__()
self.site_name = site_name
self.date_type = date_type
......@@ -65,7 +65,7 @@ class KafkaFlowAsinDetail(Templates):
# kafka相关参数(topic 按 date_type 动态:day → {site}_asin_detail_day_{yyyy_MM_dd},month → {site}_asin_detail_month_{yyyy_MM})
self.topic_name = f"{self.site_name}_asin_detail_{self.date_type}_{str(self.date_info).replace('-', '_')}"
self.batch_size = batch_size
self.batch_size_history = 100000
self.batch_size_history = 20000
self.check_path = f"/home/big_data_selection/tmp/kafka_checkpoint/{self.topic_name}_{self.consumer_type}_test" if self.test_flag == 'test' else f"/home/big_data_selection/tmp/kafka_checkpoint/{self.topic_name}_{self.consumer_type}"
self.schema = self.init_schema()
# doris相关参数:主表落 dwt 库,最新详情/父 ASIN 详情表落 selection 库
......@@ -1330,5 +1330,5 @@ if __name__ == '__main__':
test_flag = sys.argv[5]
else:
test_flag = 'normal'
handle_obj = KafkaFlowAsinDetail(site_name=site_name, date_type=date_type, date_info=date_info, consumer_type=consumer_type, test_flag=test_flag, batch_size=1200000)
handle_obj = KafkaFlowAsinDetail(site_name=site_name, date_type=date_type, date_info=date_info, consumer_type=consumer_type, test_flag=test_flag, batch_size=240000)
handle_obj.run_kafka()
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