如何无中断刷新Access Token并处理无限HTTP流?
Great question—this is a super common scenario when dealing with long-lived authenticated streams like real-time sensor data. The core challenge is refreshing your access token without leaving gaps in your stream (or minimizing interruption), since your initial stream connection relies on a token that will expire after X seconds.
Here are two robust, practical approaches to solve this, with code examples you can adapt:
Approach 1: Asyncio with aiohttp (Recommended)
Asyncio is perfect for this use case because it lets you run the stream processing and token refresh tasks concurrently without blocking each other. aiohttp is the go-to async alternative to requests and supports streaming natively.
Step-by-Step Implementation
First, install
aiohttpif you haven’t already:pip install aiohttpExample Code:
import asyncio import aiohttp # Track current token and active stream task (mutable container for safe updates) current_token = None active_stream_task = {"task": None} async def fetch_new_token(): """Replace this with your actual token refresh logic (e.g., POST to your auth endpoint).""" import time new_token = f"valid_token_{int(time.time())}" print(f"✅ Refreshed access token: {new_token[:20]}...") return new_token async def refresh_token_task(interval): """Periodically refresh the token and trigger a stream restart with the new token.""" global current_token # Initialize the first token current_token = await fetch_new_token() while True: await asyncio.sleep(interval) current_token = await fetch_new_token() # Cancel the active stream to restart with fresh token task = active_stream_task["task"] if task and not task.done(): task.cancel() print("🔄 Cancelled current stream to restart with new token.") async def stream_data(session, stream_url): """Handle the stream connection and process incoming data chunks.""" try: async with session.get( f"{stream_url}?access_token={current_token}", stream=True ) as response: response.raise_for_status() # Catch HTTP errors like 401/403 async for chunk in response.content.iter_chunked(8192): if chunk: # Process your stream data here (e.g., parse temperature readings) print(f"📥 Received chunk: {chunk[:30]}...") except asyncio.CancelledError: # Expected when token is refreshed—we'll restart the stream automatically pass except aiohttp.ClientError as e: print(f"❌ Stream error: {str(e)}. Restarting in 1 second...") await asyncio.sleep(1) raise # Trigger restart in the main loop async def main(stream_url, token_refresh_interval): async with aiohttp.ClientSession() as session: # Start the background token refresh task asyncio.create_task(refresh_token_task(token_refresh_interval)) while True: # Launch a new stream task and track it stream_task = asyncio.create_task(stream_data(session, stream_url)) active_stream_task["task"] = stream_task try: await stream_task except Exception: # Restart the stream after any errors continue if __name__ == "__main__": STREAM_URL = "https://your-stream-endpoint.com" TOKEN_REFRESH_INTERVAL = 300 # Refresh every 5 minutes (adjust as needed) asyncio.run(main(STREAM_URL, TOKEN_REFRESH_INTERVAL))
Approach 2: Threads with Requests (Synchronous)
If you prefer to stick with requests (synchronous), you can use threads to separate the stream processing and token refresh workflows. This works but requires careful handling of thread-safe signals to restart the stream.
Example Code:
import threading import time import requests from requests.exceptions import RequestException current_token = None stop_stream_event = threading.Event() TOKEN_REFRESH_INTERVAL = 300 # 5 minutes def fetch_new_token(): """Replace with your actual token refresh logic.""" global current_token import time new_token = f"valid_token_{int(time.time())}" current_token = new_token print(f"✅ Refreshed access token: {new_token[:20]}...") return new_token def token_refresh_thread(): """Background thread to handle periodic token refreshes.""" fetch_new_token() # Get initial token while True: time.sleep(TOKEN_REFRESH_INTERVAL) fetch_new_token() stop_stream_event.set() # Signal stream thread to restart def stream_thread(stream_url): """Thread to process the HTTP stream.""" session = requests.Session() while True: stop_stream_event.clear() try: print(f"🔄 Starting stream with current token...") # Add a read timeout to check for stop signals even when no chunks arrive response = session.get( f"{stream_url}?access_token={current_token}", stream=True, timeout=(5, 10) # 5s connect timeout, 10s read timeout ) response.raise_for_status() for chunk in response.iter_content(chunk_size=8192): if stop_stream_event.is_set(): break # Exit loop to restart with new token if chunk: # Process your stream data here print(f"📥 Received chunk: {chunk[:30]}...") except RequestException as e: if stop_stream_event.is_set(): print("🔄 Restarting stream with fresh token...") else: print(f"❌ Stream error: {str(e)}. Retrying in 1 second...") time.sleep(1) finally: if 'response' in locals(): response.close() def main(stream_url): # Start token refresh thread (daemon thread exits when main thread exits) threading.Thread(target=token_refresh_thread, daemon=True).start() # Start stream processing thread threading.Thread(target=stream_thread, args=(stream_url,), daemon=True).start() # Keep main thread alive while True: time.sleep(1) if __name__ == "__main__": STREAM_URL = "https://your-stream-endpoint.com" main(STREAM_URL)
Are There Any Libraries for This?
There’s no dedicated library that handles this exact scenario out of the box, but:
aiohttp(async) andrequests(sync) are the standard tools for HTTP streaming in Python.- For OAuth2-specific token refresh, libraries like
requests-oauthlibcan automate token refresh logic, but you’ll still need to integrate it with the stream restart logic shown above.
Key Considerations
- Token Refresh Reliability: Make sure your token refresh logic handles errors (like network failures or expired refresh tokens) with retries and fallback mechanisms.
- Data Loss: When restarting the stream, you might lose a small amount of data between stopping the old connection and starting the new one. If data integrity is critical, check if your stream server supports resuming from a last-known position.
- Thread Safety: When using threads, ensure shared variables (like
current_token) are accessed safely. In Python, simple assignments to strings/integers are atomic, but for complex objects, use athreading.Lock. - Graceful Shutdown: Add logic to handle program shutdown (e.g., catching
KeyboardInterrupt) to close sessions and threads cleanly.
内容的提问来源于stack exchange,提问作者moshevi

