如何将Tweepy流数据存入现有Azure CosmosDB?报错及替代方案咨询
Tweepy流数据存入Azure CosmosDB问题解决及替代方案
我想了解如何将Tweepy流数据存入已有的Azure CosmosDB数据库,当前运行脚本时出现“argument of type 'int' is not iterable”错误。同时希望知晓用Python将实时推文写入CosmosDB的其他方法,查阅微软官方文档后仍无法定位修改点。
一、解决"argument of type 'int' is not iterable"错误
你的代码报错核心原因是:lijst3是字典,for i in lijst3会遍历字典的整数类型键(tweet.id),但container.upsert_item()需要传入符合CosmosDB要求的字典对象(必须包含分区键等必要字段),直接传整数必然触发类型错误。
修复方案
- 调整数据结构:将推文信息整理为包含分区键(比如
id)的标准字典,避免用tweet.id作为字典键 - 优化客户端初始化:不要在每次收到推文时都创建
CosmosClient,提前初始化一次即可,减少资源消耗 - 简化逻辑:去掉冗余的
lijst2、lijst3变量,直接处理推文数据
修复后的关键代码:
# 提前初始化Cosmos客户端,放在类定义外 import tweepy import requests from PIL import Image from io import BytesIO import time from azure.cosmos import CosmosClient import config url = config.settings['host'] key = config.settings['master_key'] client = CosmosClient(url, credential=key) database = client.get_database_client('tweepy') container = database.get_container_client('tweepy') search_terms = ["#weer", "#regen", "#zon", "#wind", "#neerslag", "#sneeuw", "#onweer", "#droog","Schiphol"] class MyStream(tweepy.StreamingClient): def on_connect(self): print("Connected") def on_tweet(self, tweet): if tweet.referenced_tweets is None: # 整理符合CosmosDB要求的推文数据(id转字符串避免精度问题) tweet_data = { "id": str(tweet.id), "text": tweet.text, "created_at": tweet.created_at.isoformat(), "author_id": str(tweet.author_id) } # 处理图片(可选,不需要可注释) try: images = tweet.data['entities']['urls'][0]['images'] if images: tweet_data['image_url'] = images[0]['url'] response = requests.get(images[0]['url']) img = Image.open(BytesIO(response.content)) img.show() except KeyError: pass # 写入CosmosDB try: container.upsert_item(tweet_data) print(f"推文 {tweet.id} 已写入CosmosDB") except Exception as e: print(f"写入失败: {str(e)}") time.sleep(0.5) # 创建并启动流 stream = MyStream(bearer_token=bearer_token) for term in search_terms: stream.add_rules(tweepy.StreamRule(term)) stream.filter(tweet_fields=["referenced_tweets",'author_id','created_at','entities'])
二、Python实时推文写入CosmosDB的其他方法
1. Azure Functions无服务器方案
将推文流逻辑部署到Azure Functions,利用云服务自动维护运行环境,适合长期稳定的流任务:
- 创建定时器/HTTP触发的Function
- 在Function中初始化Tweepy流和Cosmos客户端
- 支持批量写入(积累N条推文后一次性写入),提升存储效率
2. Kafka + CosmosDB Connector方案
针对高吞吐量场景,先将Tweepy流数据发送到Kafka集群,再通过Azure CosmosDB的Kafka Connector自动同步数据:
- 解耦数据采集与存储环节,避免单节点瓶颈
- 支持水平扩展,应对大规模推文流量
3. 异步写入提升性能
使用Azure CosmosDB异步客户端(azure.cosmos.aio)配合asyncio,提升并发写入能力:
from azure.cosmos.aio import CosmosClient import asyncio # 异步初始化客户端 async def init_cosmos(): url = config.settings['host'] key = config.settings['master_key'] client = CosmosClient(url, credential=key) database = client.get_database_client('tweepy') return database.get_container_client('tweepy') # 异步写入方法 async def async_upsert(container, tweet_data): await container.upsert_item(tweet_data) # 在on_tweet中调用异步逻辑 def on_tweet(self, tweet): if tweet.referenced_tweets is None: tweet_data = { "id": str(tweet.id), "text": tweet.text, "created_at": tweet.created_at.isoformat() } loop = asyncio.get_event_loop() loop.run_until_complete(async_upsert(self.container, tweet_data))
内容的提问来源于stack exchange,提问作者rocketrose
相关产品推荐
相关产品推荐

