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

如何将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要求的字典对象(必须包含分区键等必要字段),直接传整数必然触发类型错误。

修复方案

  1. 调整数据结构:将推文信息整理为包含分区键(比如id)的标准字典,避免用tweet.id作为字典键
  2. 优化客户端初始化:不要在每次收到推文时都创建CosmosClient,提前初始化一次即可,减少资源消耗
  3. 简化逻辑:去掉冗余的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 03:50:23