refactor: 教学版移除全平台用户个人信息采集与持久化

- 用户 ID 转为匿名 creator_hash,昵称中间脱敏,IP/头像/主页/签名/性别不再采集
- 覆盖 xhs/weibo/bilibili/douyin/kuaishou/tieba/zhihu 7 个平台
- 删除 7 张 creator 档案 ORM 表,15 张内容/评论表新增 creator_hash 列
- B 站禁用粉丝/关注/联系人列表抓取
- 新增 tools/user_hash.py 与 4 个平台的 mock+SQLite 端到端测试

测试: pytest tests/test_no_user_info.py tests/test_weibo_no_user_info.py tests/test_douyin_no_user_info.py tests/test_kuaishou_no_user_info.py (21 passed)
This commit is contained in:
程序员阿江(Relakkes)
2026-07-01 13:09:55 +08:00
parent 8b4d8fcffa
commit 9f4f8bf768
26 changed files with 1426 additions and 822 deletions
+15 -59
View File
@@ -26,6 +26,7 @@ from typing import List
import config
from var import source_keyword_var
from tools.user_hash import anonymize_user_id, mask_nickname
from ._store_impl import *
from .bilibilli_store_media import *
@@ -62,9 +63,8 @@ async def update_bilibili_video(video_item: Dict):
"title": video_item_view.get("title", "")[:500],
"desc": video_item_view.get("desc", "")[:500],
"create_time": video_item_view.get("pubdate"),
"user_id": str(video_user_info.get("mid")),
"nickname": video_user_info.get("name"),
"avatar": video_user_info.get("face", ""),
"creator_hash": anonymize_user_id(video_user_info.get("mid")), # 创作者匿名哈希(不存原始 mid)
"nickname": mask_nickname(video_user_info.get("name")), # 用户昵称(已脱敏)
"liked_count": str(video_item_stat.get("like", "")),
"disliked_count": str(video_item_stat.get("dislike", "")),
"video_play_count": str(video_item_stat.get("view", "")),
@@ -83,22 +83,8 @@ async def update_bilibili_video(video_item: Dict):
async def update_up_info(video_item: Dict):
video_item_card_list: Dict = video_item.get("Card")
video_item_card: Dict = video_item_card_list.get("card")
saver_up_info = {
"user_id": str(video_item_card.get("mid")),
"nickname": video_item_card.get("name"),
"sex": video_item_card.get("sex"),
"sign": video_item_card.get("sign"),
"avatar": video_item_card.get("face"),
"last_modify_ts": utils.get_current_timestamp(),
"total_fans": video_item_card.get("fans"),
"total_liked": video_item_card_list.get("like_num"),
"user_rank": video_item_card.get("level_info").get("current_level"),
"is_official": video_item_card.get("official_verify").get("type"),
}
utils.logger.info(f"[store.bilibili.update_up_info] bilibili user_id:{video_item_card.get('mid')}")
await BiliStoreFactory.create_store().store_creator(creator=saver_up_info)
# 教学版:UP 主个人资料(昵称/性别/签名/头像/粉丝数等)不再落库,防骚扰。
return
async def batch_update_bilibili_video_comments(video_id: str, comments: List[Dict]):
@@ -120,11 +106,8 @@ async def update_bilibili_video_comment(video_id: str, comment_item: Dict):
"create_time": comment_item.get("ctime"),
"video_id": str(video_id),
"content": content.get("message"),
"user_id": user_info.get("mid"),
"nickname": user_info.get("uname"),
"sex": user_info.get("sex"),
"sign": user_info.get("sign"),
"avatar": user_info.get("avatar"),
"creator_hash": anonymize_user_id(user_info.get("mid")), # 创作者匿名哈希(不存原始 mid)
"nickname": mask_nickname(user_info.get("uname")), # 用户昵称(已脱敏)
"sub_comment_count": str(comment_item.get("rcount", 0)),
"like_count": like_count,
"last_modify_ts": utils.get_current_timestamp(),
@@ -149,29 +132,13 @@ async def store_video(aid, video_content, extension_file_name):
async def batch_update_bilibili_creator_fans(creator_info: Dict, fans_list: List[Dict]):
if not fans_list:
return
for fan_item in fans_list:
fan_info: Dict = {
"id": fan_item.get("mid"),
"name": fan_item.get("uname"),
"sign": fan_item.get("sign"),
"avatar": fan_item.get("face"),
}
await update_bilibili_creator_contact(creator_info=creator_info, fan_info=fan_info)
# 教学版:不再采集/存储粉丝列表(其他用户的个人信息),防骚扰。
return
async def batch_update_bilibili_creator_followings(creator_info: Dict, followings_list: List[Dict]):
if not followings_list:
return
for following_item in followings_list:
following_info: Dict = {
"id": following_item.get("mid"),
"name": following_item.get("uname"),
"sign": following_item.get("sign"),
"avatar": following_item.get("face"),
}
await update_bilibili_creator_contact(creator_info=following_info, fan_info=creator_info)
# 教学版:不再采集/存储关注列表(其他用户的个人信息),防骚扰。
return
async def batch_update_bilibili_creator_dynamics(creator_info: Dict, dynamics_list: List[Dict]):
@@ -201,26 +168,15 @@ async def batch_update_bilibili_creator_dynamics(creator_info: Dict, dynamics_li
async def update_bilibili_creator_contact(creator_info: Dict, fan_info: Dict):
save_contact_item = {
"up_id": creator_info["id"],
"fan_id": fan_info["id"],
"up_name": creator_info["name"],
"fan_name": fan_info["name"],
"up_sign": creator_info["sign"],
"fan_sign": fan_info["sign"],
"up_avatar": creator_info["avatar"],
"fan_avatar": fan_info["avatar"],
"last_modify_ts": utils.get_current_timestamp(),
}
await BiliStoreFactory.create_store().store_contact(contact_item=save_contact_item)
# 教学版:UP-粉丝关系表已移除,不再存储联系人信息。
return
async def update_bilibili_creator_dynamic(creator_info: Dict, dynamic_info: Dict):
save_dynamic_item = {
"dynamic_id": dynamic_info["dynamic_id"],
"user_id": creator_info["id"],
"user_name": creator_info["name"],
"creator_hash": anonymize_user_id(creator_info.get("id")), # 创作者匿名哈希(不存原始 ID)
"user_name": mask_nickname(creator_info.get("name")), # 用户名称(已脱敏)
"text": dynamic_info["text"],
"type": dynamic_info["type"],
"pub_ts": dynamic_info["pub_ts"],
+7 -69
View File
@@ -36,7 +36,7 @@ from sqlalchemy.orm import sessionmaker
import config
from base.base_crawler import AbstractStore
from database.db_session import get_session
from database.models import BilibiliVideoComment, BilibiliVideo, BilibiliUpInfo, BilibiliUpDynamic, BilibiliContactInfo
from database.models import BilibiliVideoComment, BilibiliVideo, BilibiliUpDynamic
from tools.async_file_writer import AsyncFileWriter
from tools import utils, words
from var import crawler_type_var
@@ -130,7 +130,6 @@ class BiliDbStoreImplement(AbstractStore):
"""
video_id = int(content_item.get("video_id"))
content_item["video_id"] = video_id
content_item["user_id"] = int(content_item.get("user_id", 0) or 0)
content_item["liked_count"] = int(content_item.get("liked_count", 0) or 0)
content_item["create_time"] = int(content_item.get("create_time", 0) or 0)
@@ -179,60 +178,12 @@ class BiliDbStoreImplement(AbstractStore):
await session.commit()
async def store_creator(self, creator: Dict):
"""
Bilibili creator DB storage implementation
Args:
creator: creator item dict
"""
creator_id = int(creator.get("user_id"))
creator["user_id"] = creator_id
creator["total_fans"] = int(creator.get("total_fans", 0) or 0)
creator["total_liked"] = int(creator.get("total_liked", 0) or 0)
creator["user_rank"] = int(creator.get("user_rank", 0) or 0)
creator["is_official"] = int(creator.get("is_official", 0) or 0)
async with get_session() as session:
result = await session.execute(select(BilibiliUpInfo).where(BilibiliUpInfo.user_id == creator_id))
creator_detail = result.scalar_one_or_none()
if not creator_detail:
creator["add_ts"] = utils.get_current_timestamp()
creator["last_modify_ts"] = utils.get_current_timestamp()
new_creator = BilibiliUpInfo(**creator)
session.add(new_creator)
else:
creator["last_modify_ts"] = utils.get_current_timestamp()
for key, value in creator.items():
setattr(creator_detail, key, value)
await session.commit()
# 教学版:UP 主个人资料不再落库
pass
async def store_contact(self, contact_item: Dict):
"""
Bilibili contact DB storage implementation
Args:
contact_item: contact item dict
"""
up_id = int(contact_item.get("up_id"))
fan_id = int(contact_item.get("fan_id"))
contact_item["up_id"] = up_id
contact_item["fan_id"] = fan_id
async with get_session() as session:
result = await session.execute(
select(BilibiliContactInfo).where(BilibiliContactInfo.up_id == up_id, BilibiliContactInfo.fan_id == fan_id)
)
contact_detail = result.scalar_one_or_none()
if not contact_detail:
contact_item["add_ts"] = utils.get_current_timestamp()
contact_item["last_modify_ts"] = utils.get_current_timestamp()
new_contact = BilibiliContactInfo(**contact_item)
session.add(new_contact)
else:
contact_item["last_modify_ts"] = utils.get_current_timestamp()
for key, value in contact_item.items():
setattr(contact_detail, key, value)
await session.commit()
# 教学版:UP-粉丝关系表已移除,不再存储联系人信息
pass
async def store_dynamic(self, dynamic_item):
"""
@@ -421,21 +372,8 @@ class BiliMongoStoreImplement(AbstractStore):
utils.logger.info(f"[BiliMongoStoreImplement.store_comment] Saved comment {comment_id} to MongoDB")
async def store_creator(self, creator_item: Dict):
"""
Store UP master information to MongoDB
Args:
creator_item: UP master data
"""
user_id = creator_item.get("user_id")
if not user_id:
return
await self.mongo_store.save_or_update(
collection_suffix="creators",
query={"user_id": user_id},
data=creator_item
)
utils.logger.info(f"[BiliMongoStoreImplement.store_creator] Saved creator {user_id} to MongoDB")
# 教学版:UP 主个人资料不再落库
pass
class BiliExcelStoreImplement:
+7 -35
View File
@@ -25,6 +25,7 @@ from typing import List
import config
from var import source_keyword_var
from tools.user_hash import anonymize_user_id, mask_nickname
from ._store_impl import *
from .douyin_store_media import *
@@ -164,18 +165,12 @@ async def update_douyin_aweme(aweme_item: Dict):
"title": aweme_item.get("desc", ""),
"desc": aweme_item.get("desc", ""),
"create_time": aweme_item.get("create_time"),
"user_id": user_info.get("uid"),
"sec_uid": user_info.get("sec_uid"),
"short_user_id": user_info.get("short_id"),
"user_unique_id": user_info.get("unique_id"),
"user_signature": user_info.get("signature"),
"nickname": user_info.get("nickname"),
"avatar": user_info.get("avatar_thumb", {}).get("url_list", [""])[0],
"creator_hash": anonymize_user_id(user_info.get("uid")), # 创作者匿名哈希(不存原始 uid)
"nickname": mask_nickname(user_info.get("nickname")), # 用户昵称(已脱敏)
"liked_count": str(interact_info.get("digg_count")),
"collected_count": str(interact_info.get("collect_count")),
"comment_count": str(interact_info.get("comment_count")),
"share_count": str(interact_info.get("share_count")),
"ip_location": aweme_item.get("ip_label", ""),
"last_modify_ts": utils.get_current_timestamp(),
"aweme_url": f"https://www.douyin.com/video/{aweme_id}",
"cover_url": _extract_content_cover_url(aweme_item),
@@ -203,20 +198,13 @@ async def update_dy_aweme_comment(aweme_id: str, comment_item: Dict):
user_info = comment_item.get("user", {})
comment_id = comment_item.get("cid")
parent_comment_id = comment_item.get("reply_id", "0")
avatar_info = (user_info.get("avatar_medium", {}) or user_info.get("avatar_300x300", {}) or user_info.get("avatar_168x168", {}) or user_info.get("avatar_thumb", {}) or {})
save_comment_item = {
"comment_id": comment_id,
"create_time": comment_item.get("create_time"),
"ip_location": comment_item.get("ip_label", ""),
"aweme_id": aweme_id,
"content": comment_item.get("text"),
"user_id": user_info.get("uid"),
"sec_uid": user_info.get("sec_uid"),
"short_user_id": user_info.get("short_id"),
"user_unique_id": user_info.get("unique_id"),
"user_signature": user_info.get("signature"),
"nickname": user_info.get("nickname"),
"avatar": avatar_info.get("url_list", [""])[0],
"creator_hash": anonymize_user_id(user_info.get("uid")), # 创作者匿名哈希(不存原始 uid)
"nickname": mask_nickname(user_info.get("nickname")), # 用户昵称(已脱敏)
"sub_comment_count": str(comment_item.get("reply_comment_total", 0)),
"like_count": (comment_item.get("digg_count") if comment_item.get("digg_count") else 0),
"last_modify_ts": utils.get_current_timestamp(),
@@ -229,24 +217,8 @@ async def update_dy_aweme_comment(aweme_id: str, comment_item: Dict):
async def save_creator(user_id: str, creator: Dict):
user_info = creator.get("user", {})
gender_map = {0: "Unknown", 1: "Male", 2: "Female"}
avatar_uri = user_info.get("avatar_300x300", {}).get("uri")
local_db_item = {
"user_id": user_id,
"nickname": user_info.get("nickname"),
"gender": gender_map.get(user_info.get("gender"), "Unknown"),
"avatar": f"https://p3-pc.douyinpic.com/img/{avatar_uri}" + r"~c5_300x300.jpeg?from=2956013662",
"desc": user_info.get("signature"),
"ip_location": user_info.get("ip_location"),
"follows": user_info.get("following_count", 0),
"fans": user_info.get("max_follower_count", 0),
"interaction": user_info.get("total_favorited", 0),
"videos_count": user_info.get("aweme_count", 0),
"last_modify_ts": utils.get_current_timestamp(),
}
utils.logger.info(f"[store.douyin.save_creator] creator:{local_db_item}")
await DouyinStoreFactory.create_store().store_creator(local_db_item)
# 教学版:创作者个人资料(昵称/性别/头像/签名/IP/粉丝数等)不再落库,防骚扰。
return
async def update_dy_aweme_image(aweme_id, pic_content, extension_file_name):
+5 -34
View File
@@ -33,7 +33,7 @@ from sqlalchemy import select
import config
from base.base_crawler import AbstractStore
from database.db_session import get_session
from database.models import DouyinAweme, DouyinAwemeComment, DyCreator
from database.models import DouyinAweme, DouyinAwemeComment
from tools import utils, words
from tools.async_file_writer import AsyncFileWriter
from var import crawler_type_var
@@ -133,24 +133,8 @@ class DouyinDbStoreImplement(AbstractStore):
await session.commit()
async def store_creator(self, creator: Dict):
"""
Douyin creator DB storage implementation
Args:
creator: creator dict
"""
user_id = creator.get("user_id")
async with get_session() as session:
result = await session.execute(select(DyCreator).where(DyCreator.user_id == user_id))
user_detail = result.scalar_one_or_none()
if not user_detail:
creator["add_ts"] = utils.get_current_timestamp()
new_creator = DyCreator(**creator)
session.add(new_creator)
else:
for key, value in creator.items():
setattr(user_detail, key, value)
await session.commit()
# 教学版:创作者个人资料不再落库
pass
class DouyinJsonStoreImplement(AbstractStore):
@@ -275,21 +259,8 @@ class DouyinMongoStoreImplement(AbstractStore):
utils.logger.info(f"[DouyinMongoStoreImplement.store_comment] Saved comment {comment_id} to MongoDB")
async def store_creator(self, creator_item: Dict):
"""
Store creator information to MongoDB
Args:
creator_item: Creator data
"""
user_id = creator_item.get("user_id")
if not user_id:
return
await self.mongo_store.save_or_update(
collection_suffix="creators",
query={"user_id": user_id},
data=creator_item
)
utils.logger.info(f"[DouyinMongoStoreImplement.store_creator] Saved creator {user_id} to MongoDB")
# 教学版:创作者个人资料不再落库
pass
class DouyinExcelStoreImplement:
+9 -25
View File
@@ -26,6 +26,7 @@ from typing import List
import config
from var import source_keyword_var
from tools.user_hash import anonymize_user_id, mask_nickname
from ._store_impl import *
@@ -63,9 +64,8 @@ async def update_kuaishou_video(video_item: Dict):
"title": photo_info.get("caption", "")[:500],
"desc": photo_info.get("caption", "")[:500],
"create_time": photo_info.get("timestamp"),
"user_id": user_info.get("id"),
"nickname": user_info.get("name"),
"avatar": user_info.get("headerUrl", ""),
"creator_hash": anonymize_user_id(user_info.get("id")), # 创作者匿名哈希(不存原始 user_id)
"nickname": mask_nickname(user_info.get("name")), # 用户昵称(已脱敏)
"liked_count": str(photo_info.get("realLikeCount")),
"viewd_count": str(photo_info.get("viewCount")),
"last_modify_ts": utils.get_current_timestamp(),
@@ -97,11 +97,10 @@ async def update_ks_video_comment(video_id: str, comment_item: Dict):
"create_time": comment_item.get("timestamp"),
"video_id": video_id,
"content": comment_item.get("content"),
# V2: author_id, Old: authorId
"user_id": comment_item.get("author_id") or comment_item.get("authorId"),
# V2: author_name, Old: authorName
"nickname": comment_item.get("author_name") or comment_item.get("authorName"),
"avatar": comment_item.get("headurl"),
# 创作者匿名哈希(不存原始 user_id):V2: author_id, Old: authorId
"creator_hash": anonymize_user_id(comment_item.get("author_id") or comment_item.get("authorId")),
# 用户昵称(已脱敏):V2: author_name, Old: authorName
"nickname": mask_nickname(comment_item.get("author_name") or comment_item.get("authorName")),
# V2: commentCount, Old: subCommentCount
"sub_comment_count": str(comment_item.get("commentCount") or comment_item.get("subCommentCount", 0)),
"last_modify_ts": utils.get_current_timestamp(),
@@ -111,20 +110,5 @@ async def update_ks_video_comment(video_id: str, comment_item: Dict):
await KuaishouStoreFactory.create_store().store_comment(comment_item=save_comment_item)
async def save_creator(user_id: str, creator: Dict):
ownerCount = creator.get('ownerCount', {})
profile = creator.get('profile', {})
local_db_item = {
'user_id': user_id,
'nickname': profile.get('user_name'),
'gender': 'Female' if profile.get('gender') == "F" else 'Male',
'avatar': profile.get('headurl'),
'desc': profile.get('user_text'),
'ip_location': "",
'follows': ownerCount.get("follow"),
'fans': ownerCount.get("fan"),
'interaction': ownerCount.get("photo_public"),
"last_modify_ts": utils.get_current_timestamp(),
}
utils.logger.info(f"[store.kuaishou.save_creator] creator:{local_db_item}")
await KuaishouStoreFactory.create_store().store_creator(local_db_item)
# 教学版:创作者个人资料(昵称/性别/头像/签名/IP/粉丝数等)不再落库,防骚扰。
return
+2 -15
View File
@@ -228,21 +228,8 @@ class KuaishouMongoStoreImplement(AbstractStore):
utils.logger.info(f"[KuaishouMongoStoreImplement.store_comment] Saved comment {comment_id} to MongoDB")
async def store_creator(self, creator_item: Dict):
"""
Store creator information to MongoDB
Args:
creator_item: Creator data
"""
user_id = creator_item.get("user_id")
if not user_id:
return
await self.mongo_store.save_or_update(
collection_suffix="creators",
query={"user_id": user_id},
data=creator_item
)
utils.logger.info(f"[KuaishouMongoStoreImplement.store_creator] Saved creator {user_id} to MongoDB")
# 教学版:创作者个人资料不再落库
pass
class KuaishouExcelStoreImplement:
+2 -4
View File
@@ -121,7 +121,5 @@ async def save_creator(user_info: TiebaCreator):
Returns:
"""
local_db_item = user_info.model_dump()
local_db_item["last_modify_ts"] = utils.get_current_timestamp()
utils.logger.info(f"[store.tieba.save_creator] creator:{local_db_item}")
await TieBaStoreFactory.create_store().store_creator(local_db_item)
# 教学版:创作者个人资料不再落库,防骚扰。
return
+5 -32
View File
@@ -35,7 +35,7 @@ from sqlalchemy.ext.asyncio import AsyncSession
import config
from base.base_crawler import AbstractStore
from database.models import TiebaNote, TiebaComment, TiebaCreator
from database.models import TiebaNote, TiebaComment
from tools import utils, words
from database.db_session import get_session
from var import crawler_type_var
@@ -137,23 +137,8 @@ class TieBaDbStoreImplement(AbstractStore):
await session.commit()
async def store_creator(self, creator: Dict):
"""
tieba content DB storage implementation
Args:
creator: creator dict
"""
user_id = creator.get("user_id")
async with get_session() as session:
stmt = select(TiebaCreator).where(TiebaCreator.user_id == user_id)
res = await session.execute(stmt)
db_creator = res.scalar_one_or_none()
if db_creator:
for key, value in creator.items():
setattr(db_creator, key, value)
else:
db_creator = TiebaCreator(**creator)
session.add(db_creator)
await session.commit()
# 教学版:创作者个人资料不再落库
pass
class TieBaJsonStoreImplement(AbstractStore):
@@ -258,20 +243,8 @@ class TieBaMongoStoreImplement(AbstractStore):
utils.logger.info(f"[TieBaMongoStoreImplement.store_comment] Saved comment {comment_id} to MongoDB")
async def store_creator(self, creator_item: Dict):
"""
Store creator information to MongoDB
Args:
creator_item: Creator data
"""
user_id = creator_item.get("user_id")
if not user_id:
return
await self.mongo_store.save_or_update(
collection_suffix="creators",
query={"user_id": user_id},
data=creator_item
)
# 教学版:创作者个人资料不再落库
pass
utils.logger.info(f"[TieBaMongoStoreImplement.store_creator] Saved creator {user_id} to MongoDB")
+18 -31
View File
@@ -25,6 +25,7 @@
import re
from typing import List
from tools.user_hash import anonymize_user_id, mask_nickname
from var import source_keyword_var
from .weibo_store_media import *
@@ -78,11 +79,13 @@ async def update_weibo_note(note_item: Dict):
if not note_item:
return
mblog: Dict = note_item.get("mblog")
user_info: Dict = mblog.get("user")
mblog: Dict = note_item.get("mblog") or {}
user_info: Dict = mblog.get("user") or {}
note_id = mblog.get("id")
content_text = mblog.get("text")
clean_text = re.sub(r"<.*?>", "", content_text)
# 教学版:原始 user_id 匿名化为 creator_hash,昵称脱敏;
# 不采集头像/主页链接/性别/IP 归属地等可定位真人的信息。
save_content_item = {
# Weibo information
"note_id": note_id,
@@ -94,14 +97,10 @@ async def update_weibo_note(note_item: Dict):
"shared_count": str(mblog.get("reposts_count", 0)),
"last_modify_ts": utils.get_current_timestamp(),
"note_url": f"https://m.weibo.cn/detail/{note_id}",
"ip_location": mblog.get("region_name", "").replace("发布于 ", ""),
# User information
"user_id": str(user_info.get("id")),
"nickname": user_info.get("screen_name", ""),
"gender": user_info.get("gender", ""),
"profile_url": user_info.get("profile_url", ""),
"avatar": user_info.get("profile_image_url", ""),
# 创作者信息(匿名化/脱敏,不含原始 user_id/avatar/gender/profile_url/ip_location)
"creator_hash": anonymize_user_id(user_info.get("id")),
"nickname": mask_nickname(user_info.get("screen_name", "")),
"source_keyword": source_keyword_var.get(),
}
utils.logger.info(f"[store.weibo.update_weibo_note] weibo note id:{note_id}, title:{save_content_item.get('content')[:24]} ...")
@@ -137,9 +136,11 @@ async def update_weibo_note_comment(note_id: str, comment_item: Dict):
if not comment_item or not note_id:
return
comment_id = str(comment_item.get("id"))
user_info: Dict = comment_item.get("user")
user_info: Dict = comment_item.get("user") or {}
content_text = comment_item.get("text")
clean_text = re.sub(r"<.*?>", "", content_text)
# 教学版:原始 user_id 匿名化为 creator_hash,昵称脱敏;
# 不采集头像/主页链接/性别/IP 归属地等可定位真人的信息。
save_comment_item = {
"comment_id": comment_id,
"create_time": utils.rfc2822_to_timestamp(comment_item.get("created_at")),
@@ -149,15 +150,11 @@ async def update_weibo_note_comment(note_id: str, comment_item: Dict):
"sub_comment_count": str(comment_item.get("total_number", 0)),
"comment_like_count": str(comment_item.get("like_count", 0)),
"last_modify_ts": utils.get_current_timestamp(),
"ip_location": comment_item.get("source", "").replace("来自", ""),
"parent_comment_id": comment_item.get("rootid", ""),
# User information
"user_id": str(user_info.get("id")),
"nickname": user_info.get("screen_name", ""),
"gender": user_info.get("gender", ""),
"profile_url": user_info.get("profile_url", ""),
"avatar": user_info.get("profile_image_url", ""),
# 创作者信息(匿名化/脱敏,不含原始 user_id/avatar/gender/profile_url/ip_location)
"creator_hash": anonymize_user_id(user_info.get("id")),
"nickname": mask_nickname(user_info.get("screen_name", "")),
}
utils.logger.info(f"[store.weibo.update_weibo_note_comment] Weibo note comment: {comment_id}, content: {save_comment_item.get('content', '')[:24]} ...")
await WeibostoreFactory.create_store().store_comment(comment_item=save_comment_item)
@@ -180,6 +177,8 @@ async def update_weibo_note_image(picid: str, pic_content, extension_file_name):
async def save_creator(user_id: str, user_info: Dict):
"""
Save creator information to local
教学版:为防骚扰不再采集/持久化创作者个人信息(昵称/性别/头像/简介/IP/粉丝数等),
此入口保留为空操作以兼容调用方。user_id 仅在调用方局部用于抓取该创作者的微博。
Args:
user_id:
user_info:
@@ -187,17 +186,5 @@ async def save_creator(user_id: str, user_info: Dict):
Returns:
"""
local_db_item = {
'user_id': user_id,
'nickname': user_info.get('screen_name'),
'gender': 'Female' if user_info.get('gender') == "f" else 'Male',
'avatar': user_info.get('avatar_hd'),
'desc': user_info.get('description'),
'ip_location': user_info.get("source", "").replace("来自", ""),
'follows': user_info.get('follow_count', ''),
'fans': user_info.get('followers_count', ''),
'tag_list': '',
"last_modify_ts": utils.get_current_timestamp(),
}
utils.logger.info(f"[store.weibo.save_creator] creator:{local_db_item}")
await WeibostoreFactory.create_store().store_creator(local_db_item)
# 教学版:创作者个人信息均不采集不持久化
return
+22 -31
View File
@@ -35,7 +35,7 @@ from sqlalchemy.ext.asyncio import AsyncSession
import config
from base.base_crawler import AbstractStore
from database.models import WeiboCreator, WeiboNote, WeiboNoteComment
from database.models import WeiboNote, WeiboNoteComment
from tools import utils, words
from tools.async_file_writer import AsyncFileWriter
from database.db_session import get_session
@@ -58,6 +58,13 @@ def calculate_number_of_files(file_store_path: str) -> int:
return 1
def _filter_model_fields(model_cls, item: Dict) -> Dict:
"""只保留目标 ORM 模型已有的列,避免把已删除/多余字段(如 avatar/gender/
profile_url/ip_location/user_id)传给 ORM 构造而报错。教学版兜底保护。"""
allowed = {col.name for col in model_cls.__table__.columns}
return {k: v for k, v in item.items() if k in allowed}
class WeiboCsvStoreImplement(AbstractStore):
def __init__(self, **kwargs):
super().__init__(**kwargs)
@@ -88,13 +95,14 @@ class WeiboCsvStoreImplement(AbstractStore):
async def store_creator(self, creator: Dict):
"""
Weibo creator CSV storage implementation
教学版:不采集/持久化创作者个人信息,空操作。
Args:
creator:
Returns:
"""
await self.writer.write_to_csv(item_type="creators", item=creator)
pass
class WeiboDbStoreImplement(AbstractStore):
@@ -108,6 +116,8 @@ class WeiboDbStoreImplement(AbstractStore):
Returns:
"""
# 教学版兜底:过滤掉已删除/多余字段,确保不会把 user_id/avatar 等传给 ORM
content_item = _filter_model_fields(WeiboNote, content_item)
note_id = int(content_item.get("note_id"))
content_item["note_id"] = note_id
async with get_session() as session:
@@ -135,6 +145,8 @@ class WeiboDbStoreImplement(AbstractStore):
Returns:
"""
# 教学版兜底:过滤掉已删除/多余字段,确保不会把 user_id/avatar 等传给 ORM
comment_item = _filter_model_fields(WeiboNoteComment, comment_item)
comment_id = int(comment_item.get("comment_id"))
comment_item["comment_id"] = comment_id
comment_item["note_id"] = int(comment_item.get("note_id", 0) or 0)
@@ -162,29 +174,14 @@ class WeiboDbStoreImplement(AbstractStore):
async def store_creator(self, creator: Dict):
"""
Weibo creator DB storage implementation
教学版:不采集/持久化创作者个人信息,空操作(WeiboCreator 表已删除)。
Args:
creator:
Returns:
"""
user_id = int(creator.get("user_id"))
creator["user_id"] = user_id
async with get_session() as session:
stmt = select(WeiboCreator).where(WeiboCreator.user_id == user_id)
res = await session.execute(stmt)
db_creator = res.scalar_one_or_none()
if db_creator:
db_creator.last_modify_ts = utils.get_current_timestamp()
for key, value in creator.items():
if hasattr(db_creator, key):
setattr(db_creator, key, value)
else:
creator["add_ts"] = utils.get_current_timestamp()
creator["last_modify_ts"] = utils.get_current_timestamp()
db_creator = WeiboCreator(**creator)
session.add(db_creator)
await session.commit()
pass
class WeiboJsonStoreImplement(AbstractStore):
@@ -217,13 +214,14 @@ class WeiboJsonStoreImplement(AbstractStore):
async def store_creator(self, creator: Dict):
"""
creator JSON storage implementation
教学版:不采集/持久化创作者个人信息,空操作。
Args:
creator:
Returns:
"""
await self.writer.write_single_item_to_json(item_type="creators", item=creator)
pass
class WeiboJsonlStoreImplement(AbstractStore):
@@ -238,7 +236,8 @@ class WeiboJsonlStoreImplement(AbstractStore):
await self.writer.write_to_jsonl(item_type="comments", item=comment_item)
async def store_creator(self, creator: Dict):
await self.writer.write_to_jsonl(item_type="creators", item=creator)
# 教学版:不采集/持久化创作者个人信息,空操作。
pass
class WeiboSqliteStoreImplement(WeiboDbStoreImplement):
@@ -291,19 +290,11 @@ class WeiboMongoStoreImplement(AbstractStore):
async def store_creator(self, creator_item: Dict):
"""
Store creator information to MongoDB
教学版:不采集/持久化创作者个人信息,空操作。
Args:
creator_item: Creator data
"""
user_id = creator_item.get("user_id")
if not user_id:
return
await self.mongo_store.save_or_update(
collection_suffix="creators",
query={"user_id": user_id},
data=creator_item
)
utils.logger.info(f"[WeiboMongoStoreImplement.store_creator] Saved creator {user_id} to MongoDB")
pass
class WeiboExcelStoreImplement:
+7 -45
View File
@@ -25,6 +25,7 @@ from typing import List
import config
from var import source_keyword_var
from tools.user_hash import anonymize_user_id, mask_nickname
from .xhs_store_media import *
from ._store_impl import *
@@ -113,14 +114,12 @@ async def update_xhs_note(note_item: Dict):
"video_url": video_url, # Note video url
"time": note_item.get("time"), # Note publish time
"last_update_time": note_item.get("last_update_time", 0), # Note last update time
"user_id": user_info.get("user_id"), # User ID
"nickname": user_info.get("nickname"), # User nickname
"avatar": user_info.get("avatar"), # User avatar
"creator_hash": anonymize_user_id(user_info.get("user_id")), # 创作者匿名哈希(不存原始 user_id)
"nickname": mask_nickname(user_info.get("nickname")), # 用户昵称(已脱敏)
"liked_count": interact_info.get("liked_count"), # Like count
"collected_count": interact_info.get("collected_count"), # Collection count
"comment_count": interact_info.get("comment_count"), # Comment count
"share_count": interact_info.get("share_count"), # Share count
"ip_location": note_item.get("ip_location", ""), # IP location
"image_list": ','.join([img.get('url', '') for img in image_list]), # Image URLs
"tag_list": ','.join([tag.get('name', '') for tag in tag_list if tag.get('type') == 'topic']), # Tags
"last_modify_ts": utils.get_current_timestamp(), # Last modification timestamp (Generated by MediaCrawler, mainly used to record the latest update time of a record in DB storage)
@@ -165,12 +164,10 @@ async def update_xhs_note_comment(note_id: str, comment_item: Dict):
local_db_item = {
"comment_id": comment_id, # Comment ID
"create_time": comment_item.get("create_time"), # Comment time
"ip_location": comment_item.get("ip_location"), # IP location
"note_id": note_id, # Note ID
"content": comment_item.get("content"), # Comment content
"user_id": user_info.get("user_id"), # User ID
"nickname": user_info.get("nickname"), # User nickname
"avatar": user_info.get("image"), # User avatar
"creator_hash": anonymize_user_id(user_info.get("user_id")), # 创作者匿名哈希(不存原始 user_id)
"nickname": mask_nickname(user_info.get("nickname")), # 用户昵称(已脱敏)
"sub_comment_count": comment_item.get("sub_comment_count", 0), # Sub-comment count
"pictures": ",".join(comment_pictures), # Comment pictures
"parent_comment_id": target_comment.get("id", 0), # Parent comment ID
@@ -191,43 +188,8 @@ async def save_creator(user_id: str, creator: Dict):
Returns:
"""
user_info = creator.get('basicInfo', {})
follows = 0
fans = 0
interaction = 0
for i in creator.get('interactions'):
if i.get('type') == 'follows':
follows = i.get('count')
elif i.get('type') == 'fans':
fans = i.get('count')
elif i.get('type') == 'interaction':
interaction = i.get('count')
def get_gender(gender):
if gender == 1:
return 'Female'
elif gender == 0:
return 'Male'
else:
return None
local_db_item = {
'user_id': user_id, # User ID
'nickname': user_info.get('nickname'), # Nickname
'gender': get_gender(user_info.get('gender')), # Gender
'avatar': user_info.get('images'), # Avatar
'desc': user_info.get('desc'), # Personal description
'ip_location': user_info.get('ipLocation'), # IP location
'follows': follows, # Following count
'fans': fans, # Fans count
'interaction': interaction, # Interaction count
'tag_list': json.dumps({tag.get('tagType'): tag.get('name')
for tag in creator.get('tags')}, ensure_ascii=False), # Tags
"last_modify_ts": utils.get_current_timestamp(), # Last modification timestamp (Generated by MediaCrawler, mainly used to record the latest update time of a record in DB storage)
}
utils.logger.info(f"[store.xhs.save_creator] creator:{local_db_item}")
await XhsStoreFactory.create_store().store_creator(local_db_item)
# 教学版:创作者个人资料(昵称/性别/头像/IP/粉丝数等)不再落库,防骚扰。
return
async def update_xhs_note_image(note_id, pic_content, extension_file_name):
+7 -65
View File
@@ -30,7 +30,7 @@ from sqlalchemy.orm import Session
from base.base_crawler import AbstractStore
from database.db_session import get_session
from database.models import XhsNote, XhsNoteComment, XhsCreator
from database.models import XhsNote, XhsNoteComment
from tools.async_file_writer import AsyncFileWriter
from tools.time_util import get_current_timestamp
@@ -137,10 +137,8 @@ class XhsDbStoreImplement(AbstractStore):
add_ts = int(get_current_timestamp())
last_modify_ts = int(get_current_timestamp())
note = XhsNote(
user_id=content_item.get("user_id"),
creator_hash=content_item.get("creator_hash"),
nickname=content_item.get("nickname"),
avatar=content_item.get("avatar"),
ip_location=content_item.get("ip_location"),
add_ts=add_ts,
last_modify_ts=last_modify_ts,
note_id=content_item.get("note_id"),
@@ -197,10 +195,8 @@ class XhsDbStoreImplement(AbstractStore):
add_ts = int(get_current_timestamp())
last_modify_ts = int(get_current_timestamp())
comment = XhsNoteComment(
user_id=comment_item.get("user_id"),
creator_hash=comment_item.get("creator_hash"),
nickname=comment_item.get("nickname"),
avatar=comment_item.get("avatar"),
ip_location=comment_item.get("ip_location"),
add_ts=add_ts,
last_modify_ts=last_modify_ts,
comment_id=comment_item.get("comment_id"),
@@ -231,54 +227,8 @@ class XhsDbStoreImplement(AbstractStore):
return result.first() is not None
async def store_creator(self, creator_item: Dict):
user_id = creator_item.get("user_id")
if not user_id:
return
async with get_session() as session:
if await self.creator_is_exist(session, user_id):
await self.update_creator(session, creator_item)
else:
await self.add_creator(session, creator_item)
async def add_creator(self, session: AsyncSession, creator_item: Dict):
add_ts = int(get_current_timestamp())
last_modify_ts = int(get_current_timestamp())
creator = XhsCreator(
user_id=creator_item.get("user_id"),
nickname=creator_item.get("nickname"),
avatar=creator_item.get("avatar"),
ip_location=creator_item.get("ip_location"),
add_ts=add_ts,
last_modify_ts=last_modify_ts,
desc=creator_item.get("desc"),
gender=creator_item.get("gender"),
follows=str(creator_item.get("follows")),
fans=str(creator_item.get("fans")),
interaction=str(creator_item.get("interaction")),
tag_list=json.dumps(creator_item.get("tag_list"))
)
session.add(creator)
async def update_creator(self, session: AsyncSession, creator_item: Dict):
user_id = creator_item.get("user_id")
last_modify_ts = int(get_current_timestamp())
update_data = {
"last_modify_ts": last_modify_ts,
"nickname": creator_item.get("nickname"),
"avatar": creator_item.get("avatar"),
"desc": creator_item.get("desc"),
"follows": str(creator_item.get("follows")),
"fans": str(creator_item.get("fans")),
"interaction": str(creator_item.get("interaction")),
"tag_list": json.dumps(creator_item.get("tag_list"))
}
stmt = update(XhsCreator).where(XhsCreator.user_id == user_id).values(**update_data)
await session.execute(stmt)
async def creator_is_exist(self, session: AsyncSession, user_id: str) -> bool:
stmt = select(XhsCreator).where(XhsCreator.user_id == user_id)
result = await session.execute(stmt)
return result.first() is not None
# 教学版:创作者个人资料不再落库
pass
async def get_all_content(self) -> List[Dict]:
async with get_session() as session:
@@ -345,16 +295,8 @@ class XhsMongoStoreImplement(AbstractStore):
Args:
creator_item: Creator data
"""
user_id = creator_item.get("user_id")
if not user_id:
return
await self.mongo_store.save_or_update(
collection_suffix="creators",
query={"user_id": user_id},
data=creator_item
)
utils.logger.info(f"[XhsMongoStoreImplement.store_creator] Saved creator {user_id} to MongoDB")
# 教学版:创作者个人资料不再落库
pass
class XhsExcelStoreImplement:
+11 -55
View File
@@ -36,7 +36,7 @@ from sqlalchemy.ext.asyncio import AsyncSession
import config
from base.base_crawler import AbstractStore
from database.db_session import get_session
from database.models import ZhihuContent, ZhihuComment, ZhihuCreator
from database.models import ZhihuContent, ZhihuComment
from tools import utils, words
from var import crawler_type_var
from tools.async_file_writer import AsyncFileWriter
@@ -85,15 +85,8 @@ class ZhihuCsvStoreImplement(AbstractStore):
await self.writer.write_to_csv(item_type="comments", item=comment_item)
async def store_creator(self, creator: Dict):
"""
Zhihu content CSV storage implementation
Args:
creator: creator dict
Returns:
"""
await self.writer.write_to_csv(item_type="creators", item=creator)
"""Creator profile is no longer persisted (teaching version: anti-harassment)."""
pass
class ZhihuDbStoreImplement(AbstractStore):
@@ -142,26 +135,8 @@ class ZhihuDbStoreImplement(AbstractStore):
await session.commit()
async def store_creator(self, creator: Dict):
"""
Zhihu content DB storage implementation
Args:
creator: creator dict
"""
user_id = creator.get("user_id")
async with get_session() as session:
stmt = select(ZhihuCreator).where(ZhihuCreator.user_id == user_id)
result = await session.execute(stmt)
existing_creator = result.scalars().first()
if existing_creator:
for key, value in creator.items():
if hasattr(existing_creator, key):
setattr(existing_creator, key, value)
else:
if "add_ts" not in creator:
creator["add_ts"] = utils.get_current_timestamp()
new_creator = ZhihuCreator(**creator)
session.add(new_creator)
await session.commit()
"""Creator profile is no longer persisted (teaching version: anti-harassment)."""
pass
class ZhihuJsonStoreImplement(AbstractStore):
@@ -192,15 +167,8 @@ class ZhihuJsonStoreImplement(AbstractStore):
await self.writer.write_single_item_to_json(item_type="comments", item=comment_item)
async def store_creator(self, creator: Dict):
"""
Zhihu content JSON storage implementation
Args:
creator: creator dict
Returns:
"""
await self.writer.write_single_item_to_json(item_type="creators", item=creator)
"""Creator profile is no longer persisted (teaching version: anti-harassment)."""
pass
class ZhihuJsonlStoreImplement(AbstractStore):
@@ -215,7 +183,8 @@ class ZhihuJsonlStoreImplement(AbstractStore):
await self.writer.write_to_jsonl(item_type="comments", item=comment_item)
async def store_creator(self, creator: Dict):
await self.writer.write_to_jsonl(item_type="creators", item=creator)
"""Creator profile is no longer persisted (teaching version: anti-harassment)."""
pass
class ZhihuSqliteStoreImplement(ZhihuDbStoreImplement):
@@ -266,21 +235,8 @@ class ZhihuMongoStoreImplement(AbstractStore):
utils.logger.info(f"[ZhihuMongoStoreImplement.store_comment] Saved comment {comment_id} to MongoDB")
async def store_creator(self, creator_item: Dict):
"""
Store creator information to MongoDB
Args:
creator_item: Creator data
"""
user_id = creator_item.get("user_id")
if not user_id:
return
await self.mongo_store.save_or_update(
collection_suffix="creators",
query={"user_id": user_id},
data=creator_item
)
utils.logger.info(f"[ZhihuMongoStoreImplement.store_creator] Saved creator {user_id} to MongoDB")
"""Creator profile is no longer persisted (teaching version: anti-harassment)."""
pass
class ZhihuExcelStoreImplement: