refactor: enhance stock sending logic to include retry mechanism and handle rate limits
This commit is contained in:
@@ -47,32 +47,35 @@ class OzonMarketplaceApi(BaseMarketplaceApi):
|
||||
|
||||
self.init_session()
|
||||
limiter = BatchLimiter()
|
||||
max_retries = 5
|
||||
while chunks:
|
||||
current_retry = 0
|
||||
chunk = chunks.pop()
|
||||
while current_retry <= max_retries:
|
||||
try:
|
||||
await limiter.acquire_ozon(self.limiter_key)
|
||||
request_data = {'stocks': chunk}
|
||||
response = await self._method('POST', '/v2/products/stocks', data=request_data)
|
||||
current_retry += 1
|
||||
response = await response.json()
|
||||
error_message = response.get('message')
|
||||
error_code = response.get('code')
|
||||
if error_message:
|
||||
if error_code == 8:
|
||||
logging.warning(f'Ozon rate limit exceeded for marketplace [{self.marketplace.id}]')
|
||||
await asyncio.sleep(1)
|
||||
continue
|
||||
else:
|
||||
logging.warning(
|
||||
f'Error occurred when sending stocks to [{self.marketplace.id}]: {error_message} ({error_code})')
|
||||
break
|
||||
else:
|
||||
break
|
||||
|
||||
async def send_stock_chunk(chunk) -> bool:
|
||||
try:
|
||||
await limiter.acquire_ozon(self.limiter_key)
|
||||
request_data = {'stocks': chunk}
|
||||
response = await self._method('POST', '/v2/products/stocks', data=request_data)
|
||||
response = await response.json()
|
||||
error_message = response.get('message')
|
||||
error_code = response.get('code')
|
||||
if error_message:
|
||||
logging.warning(
|
||||
f'Error occurred when sending stocks to [{self.marketplace.id}]: {error_message} ({error_code})')
|
||||
return False
|
||||
return True
|
||||
except Exception as e:
|
||||
logging.error(
|
||||
f'Exception occurred while sending stocks to marketplace ID [{self.marketplace.id}]: {str(e)}')
|
||||
return False
|
||||
except Exception as e:
|
||||
logging.error(
|
||||
f'Exception occurred while sending stocks to marketplace ID [{self.marketplace.id}]: {str(e)}')
|
||||
break
|
||||
|
||||
tasks = [send_stock_chunk(chunk) for chunk in chunks]
|
||||
first_request = tasks[0]
|
||||
first_response = await first_request
|
||||
if not first_response:
|
||||
logging.error(f'Skipping marketplace [{self.marketplace.id}] because first request was unsuccessful')
|
||||
await self.session.close()
|
||||
return
|
||||
|
||||
await asyncio.gather(*tasks[1:])
|
||||
await self.session.close()
|
||||
|
||||
Reference in New Issue
Block a user