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

使用Python Jupyter Notebook将带地理标签的Twitter流导入带PostGIS扩展的PostgreSQL失败并报'str' object is not callable错误

问题分析与解决方案

你遇到的两个核心问题:初期无报错但无数据写入,以及后期出现'str' object is not callable错误,我们逐一拆解解决:


1. 立即修复'str' object is not callable错误

这个报错的根源是你在dbConnect函数里把SQL字符串当成函数调用了:

# 错误写法:将字符串command当作函数执行
cur.execute(command(tweet_id, user_id, user_name, tweet, hashtags, lang, created_at, coordinates))

正确的用法是,cur.execute()第一个参数传SQL语句,第二个参数传参数元组(注意元组末尾的逗号,避免被当成单个值):

# 正确写法
cur.execute(command, (tweet_id, user_id, user_name, tweet, hashtags, lang, created_at, coordinates))

2. 解决无数据写入的深层问题

结合你的代码和数据库表结构,还有几个隐性问题导致数据无法正常入库:

2.1 表字段与插入字段名不匹配

你的PostgreSQL表定义的是username字段,但代码里插入时写的是user_name,这会导致字段不匹配的隐性错误(被try-except捕获但未提示):

-- 表结构中的字段是username
CREATE TABLE tweets (tweet_id VARCHAR PRIMARY KEY, user_id VARCHAR, username TEXT, ...);
-- 代码里误用了user_name
command = "INSERT INTO tweets (tweet_id, user_id, user_name, ...) VALUES (...)"

需要把SQL语句里的user_name改成username,和表结构保持一致。

2.2 Tweepy的on_status与on_data冲突

当你在MyStreamListener中同时重写on_status和on_data时,on_data不会被触发(因为on_status是on_data的上层封装逻辑)。你已经在on_data里处理数据入库,所以需要删除on_status方法,或者在on_status里手动调用on_data。

2.3 坐标的PostGIS存储格式错误

你的表中coordinates是GEOMETRY类型,但直接传入坐标列表无法被PostgreSQL识别。需要用PostGIS的ST_MakePoint+ST_SetSRID函数创建空间点(Twitter坐标采用WGS84坐标系,EPSG编码为4326):

# 修改SQL语句的坐标部分
command = """
INSERT INTO tweets (tweet_id, user_id, username, tweet, hashtags, lang, created_at, coordinates)
VALUES (%s,%s,%s,%s,%s,%s,%s, ST_SetSRID(ST_MakePoint(%s, %s), 4326))
"""
# 拆分坐标为经度和纬度,传入PostGIS函数
cur.execute(command, (tweet_id, user_id, user_name, tweet, json.dumps(hashtags), lang, created_at, coordinates[0], coordinates[1]))

2.4 Hashtags的格式转换

datos["entities"]["hashtags"]是一个字典列表,直接存入TEXT字段会报错,需要用json.dumps()转换成字符串:

hashtags = json.dumps(datos["entities"]["hashtags"])

2.5 时间格式的转换

Twitter返回的created_at是字符串格式(例如"Wed Oct 11 06:02:11 +0000 2023"),直接存入TIMESTAMP字段会报错,需要转换成Python的datetime对象:

from datetime import datetime
created_at = datetime.strptime(datos['created_at'], '%a %b %d %H:%M:%S %z %Y')

2.6 数据库连接的优化(可选)

每次插入都创建和关闭数据库连接效率极低,还容易出错。建议把数据库连接放在MyStreamListener的初始化方法里,只创建一次:

class MyStreamListener(tweepy.StreamListener):
    def __init__(self, time_limit=300):
        self.start_time = time.time()
        self.limit = time_limit
        # 初始化数据库连接
        try:
            self.conn = psycopg2.connect(host="localhost",database="datos_twitter",port=5433,user="xxxxxxx",password="xxxxxxx")
            self.cur = self.conn.cursor()
            print("Connected to PostgreSQL database.")
        except Exception as e:
            print(f"Database connection error: {e}")
        super(MyStreamListener, self).__init__()
    
    def on_data(self, raw_data):
        try:
            # ... 数据处理逻辑 ...
            self.cur.execute(command, (...))
            self.conn.commit()
        except Exception as e:
            print(e)
            self.conn.rollback() # 出错时回滚事务

修正后的完整代码示例

#!/usr/bin/env python
# coding: utf-8
import tweepy
import json
import psycopg2
import time
from datetime import datetime

#Insert Twitter keys
ckey = "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"
csecret = "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"
atoken = "xxxxxxxxxxxxxxxxx-xxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"
asecret = "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"

#Authorize the Twitter API
auth = tweepy.OAuthHandler(ckey, csecret)
auth.set_access_token(atoken, asecret)
api = tweepy.API(auth)

class MyStreamListener(tweepy.StreamListener):
    def __init__(self, time_limit=300):
        self.start_time = time.time()
        self.limit = time_limit
        # 初始化数据库连接
        try:
            self.conn = psycopg2.connect(host="localhost",database="datos_twitter",port=5433,user="xxxxxxx",password="xxxxxxx")
            self.cur = self.conn.cursor()
            print("Connected to PostgreSQL database.")
        except Exception as e:
            print(f"Database connection error: {e}")
        super(MyStreamListener, self).__init__()

    def on_connect(self):
        print("Connected to Twitter API.")

    def on_data(self, raw_data):
        try:
            if (time.time() - self.start_time) > self.limit:
                print("Time limit reached, stopping stream.")
                self.cur.close()
                self.conn.close()
                return False

            datos = json.loads(raw_data)
            # 只处理带有精确坐标的推文
            if datos.get("coordinates") is not None:
                tweet_id = datos['id_str']
                user_id = str(datos['user']['id']) # 转为字符串匹配表的VARCHAR类型
                user_name = datos['user']['name']
                tweet = datos['text']
                hashtags = json.dumps(datos["entities"]["hashtags"])
                lang = datos['user']['lang']
                # 转换时间格式
                created_at = datetime.strptime(datos['created_at'], '%a %b %d %H:%M:%S %z %Y')
                coordinates = datos["coordinates"]["coordinates"]

                # 构造SQL插入语句,添加冲突处理避免重复入库
                command = """
                INSERT INTO tweets (tweet_id, user_id, username, tweet, hashtags, lang, created_at, coordinates)
                VALUES (%s,%s,%s,%s,%s,%s,%s, ST_SetSRID(ST_MakePoint(%s, %s), 4326))
                ON CONFLICT (tweet_id) DO NOTHING;
                """
                self.cur.execute(command, (tweet_id, user_id, user_name, tweet, hashtags, lang, created_at, coordinates[0], coordinates[1]))
                self.conn.commit()
                print(f"Inserted tweet: {tweet_id}")
        except Exception as e:
            print(f"Error processing tweet: {e}")
            self.conn.rollback()

    def on_error(self, status_code):
        print(f"Error code: {status_code}")
        if status_code == 420:
            self.cur.close()
            self.conn.close()
            return False

#Streaming of tweets
myStreamListener = MyStreamListener(time_limit=300)
myStream = tweepy.Stream(auth=api.auth, listener=myStreamListener, tweet_mode="extended")
#Filtering of tweets by spatial box and keywords
myStream.filter(locations=[-10.78,34.15, 5.95,44.04], track=['Madrid', 'madrid'])

额外检查点

  • 确保PostGIS扩展已启用:执行CREATE EXTENSION IF NOT EXISTS postgis;
  • 确认PostgreSQL端口配置正确(你用的是5433,默认是5432)
  • 验证数据库用户拥有tweets表的插入权限
  • 若仍无数据,可暂时移除坐标过滤逻辑,先测试无坐标推文的入库情况,排除过滤条件过严的问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 22:59:05