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

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 blocking requests calls 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: asyncpg connection pool ensures database inserts don’t block stream processing, critical for handling high volumes of concurrent data.
  • Concurrency: asyncio.gather runs all device streams in parallel, maximizing efficiency compared to sequential processing.

Critical Next Steps:

  1. Update Parsing Logic: Modify the parse_stream_line method to match your device’s actual HTTP stream format. The example is just a placeholder.
  2. Adjust Database Schema: Ensure your PostgreSQL table (device_mtr_data in the example) matches the columns you’re inserting.
  3. 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).
  4. Tune Timeouts: Adjust connection timeouts based on your devices’ response characteristics.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 04:22:27