feat(新功能): batch generate 增加接口说明,design batch队列修改

fix(修复bug):
docs(文档变更):
refactor(重构):
test(增加测试):
This commit is contained in:
zchengrong
2025-04-22 13:59:59 +08:00
parent af8ed730cc
commit 96002eb7f2
5 changed files with 19 additions and 18 deletions

View File

@@ -102,12 +102,12 @@ def batch_generate_product(batch_request_data):
if DEBUG is False:
if i + 1 < batch_size:
publish_status(tasks_id, f"{i + 1}/{batch_size}", image_url)
logger.info(f" [x] {tasks_id}tasks_id *** progress{i + 1}/{batch_size} *** image_url{image_url}")
print(f" [x] {tasks_id}tasks_id *** progress{i + 1}/{batch_size} *** image_url{image_url}")
logger.info(f" [x]Queue : {BATCH_GPI_RABBITMQ_QUEUES} | tasks_id{tasks_id} | progress{i + 1}/{batch_size} | image_url{image_url}")
# print(f" [x]Queue : {BATCH_GPI_RABBITMQ_QUEUES} | tasks_id{tasks_id} | progress{i + 1}/{batch_size} | image_url{image_url}")
else:
publish_status(tasks_id, f"OK", image_url_list)
logger.info(f" [x] {tasks_id}tasks_id *** progressOK *** image_url{image_url_list}")
print(f" [x] {tasks_id}tasks_id *** progressOK *** image_url{image_url_list}")
logger.info(f" [x]Queue : {BATCH_GPI_RABBITMQ_QUEUES} | tasks_id{tasks_id} | progressOK | image_url{image_url}")
# print(f" [x]Queue : {BATCH_GPI_RABBITMQ_QUEUES} | tasks_id{tasks_id} | progressOK | image_url{image_url}")
def pre_processing_image(image_url):

View File

@@ -127,12 +127,12 @@ def batch_generate_relight(batch_request_data):
if DEBUG is False:
if i + 1 < batch_size:
publish_status(tasks_id, f"{i + 1}/{batch_size}", image_url)
logger.info(f" [x] {tasks_id}tasks_id *** progress{i + 1}/{batch_size} *** image_url{image_url}")
print(f" [x] {tasks_id}tasks_id *** progress{i + 1}/{batch_size} *** image_url{image_url}")
logger.info(f" [x]Queue : {BATCH_GRI_RABBITMQ_QUEUES} | tasks_id{tasks_id} | progress{i + 1}/{batch_size} | image_url{image_url}")
# print(f" [x]Queue : {BATCH_GRI_RABBITMQ_QUEUES} | tasks_id{tasks_id} | progress{i + 1}/{batch_size} | image_url{image_url}")
else:
publish_status(tasks_id, f"OK", image_url_list)
logger.info(f" [x] {tasks_id}tasks_id *** progressOK *** image_url{image_url_list}")
print(f" [x] {tasks_id}tasks_id *** progressOK *** image_url{image_url_list}")
logger.info(f" [x]Queue : {BATCH_GRI_RABBITMQ_QUEUES} | tasks_id{tasks_id} | progressOK | image_url{image_url}")
# print(f" [x]Queue : {BATCH_GRI_RABBITMQ_QUEUES} | tasks_id{tasks_id} | progressOK | image_url{image_url}")
def publish_status(task_id, progress, result):

View File

@@ -29,7 +29,7 @@ from app.service.utils.oss_client import oss_get_image
minio_client = Minio(MINIO_URL, access_key=MINIO_ACCESS, secret_key=MINIO_SECRET, secure=MINIO_SECURE)
logger = logging.getLogger()
celery_app = Celery('tasks', broker=f'amqp://rabbit:123456@18.167.251.121:5672//', backend='rpc://', BROKER_CONNECTION_RETRY_ON_STARTUP=True)
celery_app = Celery('post_transform_tasks', broker=f'amqp://rabbit:123456@18.167.251.121:5672//', backend='rpc://', BROKER_CONNECTION_RETRY_ON_STARTUP=True)
celery_app.conf.task_default_queue = 'queue_post_transform'
celery_app.conf.worker_log_format = '%(asctime)s %(filename)s [line:%(lineno)d] %(levelname)s %(message)s'
celery_app.conf.worker_hijack_root_logger = False
@@ -144,21 +144,21 @@ def batch_generate_pose_transform(batch_request_data):
if DEBUG is False:
if i + 1 < batch_size:
publish_status(tasks_id, f"{i + 1}/{batch_size}", data)
logger.info(f" [x] {tasks_id}tasks_id *** progress{i + 1}/{batch_size} *** image_url{data}")
print(f" [x] {tasks_id}tasks_id *** progress{i + 1}/{batch_size} *** image_url{data}")
logger.info(f" [x]Queue : {BATCH_PS_RABBITMQ_QUEUES} | tasks_id{tasks_id} | progress{i + 1}/{batch_size} | image_url{image_url}")
# print(f" [x]Queue : {BATCH_GRI_RABBITMQ_QUEUES} | tasks_id{tasks_id} | progress{i + 1}/{batch_size} | image_url{image_url}")
else:
publish_status(tasks_id, f"OK", result_url_list)
logger.info(f" [x] {tasks_id}tasks_id *** progressOK *** image_url{result_url_list}")
print(f" [x] {tasks_id}tasks_id *** progressOK *** image_url{result_url_list}")
logger.info(f" [x]Queue : {BATCH_PS_RABBITMQ_QUEUES} | tasks_id{tasks_id} | progressOK | image_url{image_url}")
# print(f" [x]Queue : {BATCH_PS_RABBITMQ_QUEUES} | tasks_id{tasks_id} | progressOK | image_url{image_url}")
def publish_status(task_id, progress, result):
connection = pika.BlockingConnection(pika.ConnectionParameters(**RABBITMQ_PARAMS))
channel = connection.channel()
channel.queue_declare(queue=BATCH_GRI_RABBITMQ_QUEUES, durable=True)
channel.queue_declare(queue=BATCH_PS_RABBITMQ_QUEUES, durable=True)
message = {'task_id': task_id, 'progress': progress, "result": result}
channel.basic_publish(exchange='',
routing_key=BATCH_GRI_RABBITMQ_QUEUES,
routing_key=BATCH_PS_RABBITMQ_QUEUES,
body=json.dumps(message),
properties=pika.BasicProperties(
delivery_mode=2,