# -*- coding: utf-8 -*- """检测模块:并发检测 IPTV 源可用性""" import requests from concurrent.futures import ThreadPoolExecutor, as_completed from config import CHECK_TIMEOUT, MAX_WORKERS, HEADERS def check_channel(channel): """检测单个频道是否可用 返回 (channel, is_ok) 元组 """ url = channel.get("url", "") if not url: return channel, False try: # 使用流式请求,只读取少量数据 resp = requests.get( url, headers=HEADERS, timeout=CHECK_TIMEOUT, stream=True, allow_redirects=True, ) # 状态码检查 if resp.status_code != 200: return channel, False # 读取前 1KB 数据验证内容 content = b"" for chunk in resp.iter_content(chunk_size=1024): content += chunk if len(content) >= 1024: break # 检查是否为有效的流内容(m3u8 或 ts 流) if content: text = content[:200].decode("utf-8", errors="ignore") # m3u8 文件或以 #EXTM3U 开头 if "#EXTM3U" in text or "#EXTINF" in text: return channel, True # ts 流通常以 0x47 开头 if content[:1] == b"\x47": return channel, True # 其他情况认为可用(有些流返回的是二进制数据) if len(content) > 100: return channel, True return channel, False except Exception: return channel, False def check_all(channels): """并发检测所有频道""" print(f"[检测] 开始检测 {len(channels)} 个频道,并发数 {MAX_WORKERS}") valid = [] total = len(channels) done = 0 with ThreadPoolExecutor(max_workers=MAX_WORKERS) as executor: futures = {executor.submit(check_channel, ch): ch for ch in channels} for future in as_completed(futures): done += 1 channel, ok = future.result() if ok: valid.append(channel) if done % 50 == 0 or done == total: print(f"[检测进度] {done}/{total},可用 {len(valid)}") print(f"[检测完成] 可用频道 {len(valid)}/{total}") return valid if __name__ == "__main__": from collector import collect_all channels = collect_all() valid = check_all(channels) print(f"最终可用:{len(valid)}")