fix: 增加 busy_timeout、最小化事务块、增加每批处理 page 量 (#420)

This commit is contained in:
ᴀᴍᴛᴏᴀᴇʀ
2025-08-06 14:08:07 +08:00
committed by GitHub
parent dbcb1fa78b
commit 4f780faf64
4 changed files with 60 additions and 54 deletions
+18 -3
View File
@@ -1,6 +1,10 @@
use std::time::Duration;
use anyhow::{Context, Result}; use anyhow::{Context, Result};
use bili_sync_migration::{Migrator, MigratorTrait}; use bili_sync_migration::{Migrator, MigratorTrait};
use sea_orm::{ConnectOptions, Database, DatabaseConnection}; use sea_orm::sqlx::sqlite::SqliteConnectOptions;
use sea_orm::sqlx::{ConnectOptions as SqlxConnectOptions, Sqlite};
use sea_orm::{ConnectOptions, Database, DatabaseConnection, SqlxSqliteConnector};
use crate::config::CONFIG_DIR; use crate::config::CONFIG_DIR;
@@ -13,8 +17,19 @@ async fn database_connection() -> Result<DatabaseConnection> {
option option
.max_connections(100) .max_connections(100)
.min_connections(5) .min_connections(5)
.acquire_timeout(std::time::Duration::from_secs(90)); .acquire_timeout(Duration::from_secs(90));
Ok(Database::connect(option).await?) let connect_option = option
.get_url()
.parse::<SqliteConnectOptions>()
.context("Failed to parse database URL")?
.disable_statement_logging()
.busy_timeout(Duration::from_secs(90));
Ok(SqlxSqliteConnector::from_sqlx_sqlite_pool(
option
.sqlx_pool_options::<Sqlite>()
.connect_with(connect_option)
.await?,
))
} }
async fn migrate_database() -> Result<()> { async fn migrate_database() -> Result<()> {
+2 -5
View File
@@ -152,10 +152,7 @@ impl VideoInfo {
} }
impl PageInfo { impl PageInfo {
pub fn into_active_model( pub fn into_active_model(self, video_model_id: i32) -> bili_sync_entity::page::ActiveModel {
self,
video_model: &bili_sync_entity::video::Model,
) -> bili_sync_entity::page::ActiveModel {
let (width, height) = match &self.dimension { let (width, height) = match &self.dimension {
Some(d) => { Some(d) => {
if d.rotate == 0 { if d.rotate == 0 {
@@ -167,7 +164,7 @@ impl PageInfo {
None => (None, None), None => (None, None),
}; };
bili_sync_entity::page::ActiveModel { bili_sync_entity::page::ActiveModel {
video_id: Set(video_model.id), video_id: Set(video_model_id),
cid: Set(self.cid), cid: Set(self.cid),
pid: Set(self.page), pid: Set(self.page),
name: Set(self.name), name: Set(self.name),
+3 -11
View File
@@ -6,7 +6,7 @@ use sea_orm::sea_query::{OnConflict, SimpleExpr};
use sea_orm::{DatabaseTransaction, TransactionTrait}; use sea_orm::{DatabaseTransaction, TransactionTrait};
use crate::adapter::{VideoSource, VideoSourceEnum}; use crate::adapter::{VideoSource, VideoSourceEnum};
use crate::bilibili::{PageInfo, VideoInfo}; use crate::bilibili::VideoInfo;
use crate::config::{Config, LegacyConfig}; use crate::config::{Config, LegacyConfig};
use crate::utils::status::STATUS_COMPLETED; use crate::utils::status::STATUS_COMPLETED;
@@ -73,16 +73,8 @@ pub async fn create_videos(
} }
/// 尝试创建 Page Model,如果发生冲突则忽略 /// 尝试创建 Page Model,如果发生冲突则忽略
pub async fn create_pages( pub async fn create_pages(pages_model: Vec<page::ActiveModel>, connection: &DatabaseTransaction) -> Result<()> {
pages_info: Vec<PageInfo>, for page_chunk in pages_model.chunks(200) {
video_model: &bili_sync_entity::video::Model,
connection: &DatabaseTransaction,
) -> Result<()> {
let page_models = pages_info
.into_iter()
.map(|p| p.into_active_model(video_model))
.collect::<Vec<page::ActiveModel>>();
for page_chunk in page_models.chunks(50) {
page::Entity::insert_many(page_chunk.to_vec()) page::Entity::insert_many(page_chunk.to_vec())
.on_conflict( .on_conflict(
OnConflict::columns([page::Column::VideoId, page::Column::Pid]) OnConflict::columns([page::Column::VideoId, page::Column::Pid])
+10 -8
View File
@@ -106,8 +106,7 @@ pub async fn fetch_video_details(
let semaphore_ref = &semaphore; let semaphore_ref = &semaphore;
let tasks = videos_model let tasks = videos_model
.into_iter() .into_iter()
.map(|video_model| { .map(|video_model| async move {
async move {
let _permit = semaphore_ref.acquire().await.context("acquire semaphore failed")?; let _permit = semaphore_ref.acquire().await.context("acquire semaphore failed")?;
let video = Video::new(bili_client, video_model.bvid.clone()); let video = Video::new(bili_client, video_model.bvid.clone());
let info: Result<_> = async { Ok((video.get_tags().await?, video.get_view_info().await?)) }.await; let info: Result<_> = async { Ok((video.get_tags().await?, video.get_view_info().await?)) }.await;
@@ -127,21 +126,24 @@ pub async fn fetch_video_details(
let VideoInfo::Detail { pages, .. } = &mut view_info else { let VideoInfo::Detail { pages, .. } = &mut view_info else {
unreachable!() unreachable!()
}; };
// 构造 page model
let pages = std::mem::take(pages); let pages = std::mem::take(pages);
let pages_len = pages.len(); let pages = pages
let txn = connection.begin().await?; .into_iter()
// 将分页信息写入数据库 .map(|p| p.into_active_model(video_model.id))
create_pages(pages, &video_model, &txn).await?; .collect::<Vec<page::ActiveModel>>();
// 更新 video model 的各项有关属性
let mut video_active_model = view_info.into_detail_model(video_model); let mut video_active_model = view_info.into_detail_model(video_model);
video_source.set_relation_id(&mut video_active_model); video_source.set_relation_id(&mut video_active_model);
video_active_model.single_page = Set(Some(pages_len == 1)); video_active_model.single_page = Set(Some(pages.len() == 1));
video_active_model.tags = Set(Some(serde_json::to_value(tags)?)); video_active_model.tags = Set(Some(serde_json::to_value(tags)?));
let txn = connection.begin().await?;
create_pages(pages, &txn).await?;
video_active_model.save(&txn).await?; video_active_model.save(&txn).await?;
txn.commit().await?; txn.commit().await?;
} }
}; };
Ok::<_, anyhow::Error>(()) Ok::<_, anyhow::Error>(())
}
}) })
.collect::<FuturesUnordered<_>>(); .collect::<FuturesUnordered<_>>();
tasks.try_collect::<Vec<_>>().await?; tasks.try_collect::<Vec<_>>().await?;