Python异步HTTP请求:设备无限HTTP页面解析入库技术咨询
Convert to Async HTTP Requests with Data Parsing & Database Storage
Let’s refactor your synchronous code to handle concurrent device streams efficiently using async Python. Here’s a complete implementation that uses aiohttp for async HTTP calls and asyncpg for non-blocking database operations:
import asyncio import asyncpg import aiohttp from aiohttp import ClientTimeout class NodeHttpMtr: def __init__(self, ip, nsrc, ndst): self.ip = ip self.nsrc = nsrc self.ndst = ndst self.url = f"http://{self.ip}/nph-cgi_mtr?duration=-1&interval=0" self.response = None self.reader = None async def connect(self, session): """Establish async connection to the device's streaming endpoint""" try: # Set timeout only for initial connection, not the infinite stream connect_timeout = ClientTimeout(total=10) self.response = await session.get( self.url, timeout=connect_timeout, timeout=None # Disable timeout for the ongoing stream ) self.response.raise_for_status() self.reader = self.response.content except Exception as e: print(f"Failed to connect to {self.ip}: {str(e)}") raise async def __aiter__(self): """Async iterator setup""" return self async def __anext__(self): """Read and parse next line from the infinite stream""" try: line = await self.reader.readline() if not line: raise StopAsyncIteration # Parse the line - replace this with your actual data parsing logic parsed_data = self.parse_stream_line(line.decode('utf-8').strip()) if parsed_data: return parsed_data else: # Skip empty or invalid lines and fetch next return await self.__anext__() except Exception as e: print(f"Error reading stream from {self.ip}: {str(e)}") raise StopAsyncIteration def parse_stream_line(self, line): """Convert raw stream line into structured data for DB storage""" # Example parsing logic - adjust to match your device's output format if line.startswith("MTR_DATA:"): data_parts = line.split("|") return { "source": self.nsrc, "destination": self.ndst, "device_ip": self.ip, "latency": float(data_parts[1]), "packet_loss": int(data_parts[2]), "timestamp": data_parts[3] } return None async def insert_to_database(db_pool, data): """Async function to store parsed data into PostgreSQL""" async with db_pool.acquire() as conn: try: await conn.execute( """ INSERT INTO device_mtr_data (source, destination, device_ip, latency, packet_loss, timestamp) VALUES ($1, $2, $3, $4, $5, $6) """, data["source"], data["destination"], data["device_ip"], data["latency"], data["packet_loss"], data["timestamp"] ) except Exception as e: print(f"DB insert failed for {data['device_ip']}: {str(e)}") async def process_single_device(ip, nsrc, ndst, db_pool, session): """Handle full lifecycle for one device: connect, stream, parse, store""" mtr_stream = NodeHttpMtr(ip, nsrc, ndst) try: await mtr_stream.connect(session) async for parsed_data in mtr_stream: await insert_to_database(db_pool, parsed_data) except Exception as e: print(f"Processing error for {ip}: {str(e)}") finally: if mtr_stream.response: await mtr_stream.response.release() async def main(): # Initialize PostgreSQL connection pool db_pool = await asyncpg.create_pool( user="your_db_user", password="your_db_password", database="your_db_name", host="localhost" ) # List of your devices - add all your devices here devices = [ {"ip": "192.168.1.10", "nsrc": "src_office", "ndst": "dst_datacenter"}, {"ip": "192.168.1.11", "nsrc": "src_warehouse", "ndst": "dst_datacenter"}, # Add more devices as needed ] # Create a single aiohttp session for all device connections async with aiohttp.ClientSession() as session: # Run all device processing tasks concurrently processing_tasks = [ process_single_device(dev["ip"], dev["nsrc"], dev["ndst"], db_pool, session) for dev in devices ] await asyncio.gather(*processing_tasks) # Cleanup database pool await db_pool.close() if __name__ == "__main__": asyncio.run(main())
Key Improvements & Explanations:
- Async HTTP with
aiohttp: Replaces blockingrequestscalls with async sessions, allowing concurrent stream handling across multiple devices. - Async Iterator: Uses
__aiter__and__anext__(async counterparts to your original iterator) to read the infinite stream without blocking other tasks. - Non-blocking Database:
asyncpgconnection pool ensures database inserts don’t block stream processing, critical for handling high volumes of concurrent data. - Concurrency:
asyncio.gatherruns all device streams in parallel, maximizing efficiency compared to sequential processing.
Critical Next Steps:
- Update Parsing Logic: Modify the
parse_stream_linemethod to match your device’s actual HTTP stream format. The example is just a placeholder. - Adjust Database Schema: Ensure your PostgreSQL table (
device_mtr_datain the example) matches the columns you’re inserting. - Add Reconnection Logic: For production, add retry logic if a device drops the connection (e.g., wrap the stream reading in a loop with backoff).
- Tune Timeouts: Adjust connection timeouts based on your devices’ response characteristics.
内容的提问来源于stack exchange,提问作者tiler
相关产品推荐
相关产品推荐

