Asyncio环境下新老系统并行执行与结果留存需求的四种解决方案利弊咨询
First, let’s recap your scenario: you’re replacing a legacy system with a Step Functions state machine, and now need to compare outputs by running both systems in parallel. The problem arises because boto3’s synchronous calls don’t integrate with asyncio’s event loop, causing async.wait to return both tasks immediately instead of waiting for the first completion.
Let’s go through each solution’s pros, cons, and hidden pitfalls:
1. Don’t cancel pending tasks—wait for them to complete
- Pros: Dead simple to implement. You can add your result parsing/persistence logic directly in
legacy_system()without touching the async core. No changes needed to how you call boto3. - Cons:
- Lambda Concurrency Impact: This is the big one. If your new system handles, say, 250 QPS and each legacy Lambda runs for 3-4 seconds, you’ll immediately hit AWS Lambda’s default concurrent execution limit (1000). Any additional requests will get throttled until existing invocations finish. If your QPS is higher, this becomes a critical bottleneck.
- Cost: You’ll pay for full execution time of every legacy Lambda invocation, even when the new system returns first. Over time, this adds up.
- Downstream Load: If your legacy Lambda calls other services/databases, you’ll be forcing unnecessary load on those systems for every request.
- Hidden Pitfall: Make sure unhandled exceptions in the legacy Lambda don’t bubble up to your main async function. Even if you don’t care about the result, a crash could take down your main process if not caught properly.
2. Use async executors to integrate boto3 with the event loop
Pros: This fixes the root cause of your problem. By wrapping blocking boto3 calls in async-compatible execution,
async.waitwill behave as expected—waiting for the first task to complete. It preserves your original asyncio design, keeps your main loop non-blocking, and doesn’t waste resources on unnecessary waits.Cons: Requires a bit more code setup (but it’s straightforward once you know how). You’ll need to either use a thread pool or switch to an async AWS SDK.
Hidden Pitfalls:
- If using thread pools, you need to tune the
max_workersvalue to match your QPS and Lambda execution time. Too few workers will cause task queuing; too many can waste resources. - Boto3 clients are thread-safe, but avoid modifying client configurations across threads to prevent unexpected behavior.
- If using thread pools, you need to tune the
Implementation Fix:
The easiest way is to useaioboto3(the async AWS SDK) instead of regular boto3—it’s designed to work natively with asyncio, no thread pools needed. Here’s how to adjust your code:First, install it:
pip install aioboto3Then update your async functions:
import aioboto3 import asyncio async def new_system(): async with aioboto3.client("stepfunctions") as sfn_client: return await sfn_client.start_sync_execution( stateMachineArn="your-express-state-machine-arn", input='{"your": "input"}' ) async def legacy_system(): async with aioboto3.client("lambda") as lambda_client: response = await lambda_client.invoke( FunctionName="your-legacy-lambda-arn", InvocationType="RequestResponse", Payload='{"your": "request-payload"}' ) # Parse and persist the result here (e.g., save to S3/DynamoDB) payload = await response["Payload"].read() return payload async def wait_first(): new_task = asyncio.create_task(new_system()) legacy_task = asyncio.create_task(legacy_system()) done, pending = await asyncio.wait( [new_task, legacy_task], return_when=asyncio.FIRST_COMPLETED ) for task in done: try: result = await task # Cancel pending tasks (note: Lambda will still finish executing, we just stop waiting) for pending_task in pending: pending_task.cancel() try: await pending_task except asyncio.CancelledError: pass return result except Exception: # New system failed, fall back to legacy await asyncio.gather(*pending) for pending_task in pending: return await pending_taskIf you can’t use
aioboto3, you can wrap synchronous boto3 calls in a thread pool:import asyncio import boto3 from concurrent.futures import ThreadPoolExecutor # Configure thread pool size based on your QPS needs executor = ThreadPoolExecutor(max_workers=10) def sync_new_system(): sfn_client = boto3.client("stepfunctions") return sfn_client.start_sync_execution( stateMachineArn="your-express-state-machine-arn", input='{"your": "input"}' ) def sync_legacy_system(): lambda_client = boto3.client("lambda") response = lambda_client.invoke( FunctionName="your-legacy-lambda-arn", InvocationType="RequestResponse", Payload='{"your": "request-payload"}' ) payload = response["Payload"].read() # Add persistence logic here return payload async def new_system(): loop = asyncio.get_event_loop() return await loop.run_in_executor(executor, sync_new_system) async def legacy_system(): loop = asyncio.get_event_loop() return await loop.run_in_executor(executor, sync_legacy_system) # wait_first function remains the same as above
3. Use Lambda async invocation (Event type) + polling
- Pros: Offloads the wait to Lambda’s async execution, so your main function doesn’t block on the legacy call.
- Cons:
- Polling is an anti-pattern for event-driven architectures. It adds unnecessary API calls (and cost) to check for Lambda results.
- Tuning polling intervals is tricky: too frequent and you waste resources; too slow and you delay result persistence.
- Lambda’s async invocation results are only retrievable for 24 hours via
GetFunctionInvocation—if your legacy Lambda runs longer than that, you’ll lose the result.
- Hidden Pitfall: You’ll need to handle retries for when the Lambda is still running, and manage state for in-flight requests to avoid duplicate polling. This adds significant complexity to your codebase.
4. Add await asyncio.sleep()
- Pros: The quickest hack to implement.
- Cons: Extremely fragile. The sleep time is guesswork—if your new system’s response time varies (e.g., under load), you’ll either still get simultaneous task returns (sleep too short) or introduce unnecessary latency (sleep too long). It also blocks the async event loop, reducing your main function’s concurrency capacity.
- Hidden Pitfall: This is a band-aid, not a fix. It will break as soon as your system’s performance characteristics change, and it’s confusing for future developers to maintain.
Final Recommendation
Option 2 is the clear winner here. Using aioboto3 is the cleanest, most performant approach since it’s natively async and avoids thread pool overhead. If you can’t use aioboto3, the thread pool method works reliably as long as you tune the worker count appropriately.
内容的提问来源于stack exchange,提问作者lynkfox

