使用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

