You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何无中断刷新Access Token并处理无限HTTP流?

Handling Infinite HTTP Streams with Periodic Access Token Refresh

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:

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

  1. First, install aiohttp if you haven’t already:

    pip install aiohttp
    
  2. Example 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) and requests (sync) are the standard tools for HTTP streaming in Python.
  • For OAuth2-specific token refresh, libraries like requests-oauthlib can automate token refresh logic, but you’ll still need to integrate it with the stream restart logic shown above.

Key Considerations

  1. Token Refresh Reliability: Make sure your token refresh logic handles errors (like network failures or expired refresh tokens) with retries and fallback mechanisms.
  2. 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.
  3. 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 a threading.Lock.
  4. Graceful Shutdown: Add logic to handle program shutdown (e.g., catching KeyboardInterrupt) to close sessions and threads cleanly.

内容的提问来源于stack exchange,提问作者moshevi

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.06 23:32:35