Commit 76d72790 by chenyuanjie

fix

parent fecee2ab
......@@ -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=20000):
def __init__(self, site_name='us', date_type="month", date_info='2026-03', consumer_type='history', test_flag='test', batch_size=1200000):
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 = 20000
self.batch_size_history = 100000
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 库
......@@ -1308,5 +1308,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=20000)
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.run_kafka()
......@@ -132,6 +132,7 @@ class Templates(object):
.option("failOnDataLoss", "false")
# 断点与topic最新offset差距过大时,避免单个触发批次读取过多数据导致阻塞/失败风险,按批次上限分批追平
max_offsets_per_trigger = getattr(self, 'batch_size', None)
print(f"[create_kafka_df_object] maxOffsetsPerTrigger 将设置为: {max_offsets_per_trigger}")
if max_offsets_per_trigger:
kafka_stream_reader = kafka_stream_reader.option("maxOffsetsPerTrigger", max_offsets_per_trigger)
kafka_df = kafka_stream_reader.load() \
......
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