Database optimizations.

This commit is contained in:
Jordan Petridis
2017-10-23 10:47:54 +03:00
parent 9beea21a4f
commit 0e5d976514
6 changed files with 58 additions and 36 deletions
+32 -23
View File
@@ -98,7 +98,7 @@ pub fn index_loop(db: &Arc<Mutex<SqliteConnection>>, force: bool) -> Result<()>
pub fn complete_index_from_source(
req: &mut reqwest::Response,
source: &Source,
mutex: &Arc<Mutex<SqliteConnection>>,
db: &Arc<Mutex<SqliteConnection>>,
) -> Result<()> {
use std::io::Read;
use std::str::FromStr;
@@ -107,57 +107,66 @@ pub fn complete_index_from_source(
req.read_to_string(&mut buf)?;
let chan = rss::Channel::from_str(&buf)?;
complete_index(mutex, &chan, source)?;
complete_index(&db, &chan, source)?;
Ok(())
}
fn complete_index(
connection: &Arc<Mutex<SqliteConnection>>,
db: &Arc<Mutex<SqliteConnection>>,
chan: &rss::Channel,
parent: &Source,
) -> Result<()> {
let pd = {
let db = connection.lock().unwrap();
index_channel(&db, chan, parent)?
let conn = db.lock().unwrap();
index_channel(&conn, chan, parent)?
};
index_channel_items(connection, chan.items(), &pd);
index_channel_items(db, chan.items(), &pd);
Ok(())
}
fn index_channel(db: &SqliteConnection, chan: &rss::Channel, parent: &Source) -> Result<Podcast> {
fn index_channel(con: &SqliteConnection, chan: &rss::Channel, parent: &Source) -> Result<Podcast> {
let pd = feedparser::parse_podcast(chan, parent.id());
// Convert NewPodcast to Podcast
let pd = insert_return_podcast(db, &pd)?;
let pd = insert_return_podcast(con, &pd)?;
Ok(pd)
}
fn index_channel_items(connection: &Arc<Mutex<SqliteConnection>>, it: &[rss::Item], pd: &Podcast) {
it.par_iter()
fn index_channel_items(db: &Arc<Mutex<SqliteConnection>>, it: &[rss::Item], pd: &Podcast) {
let episodes: Vec<_> = it.par_iter()
.map(|x| feedparser::parse_episode(x, pd.id()))
.for_each(|x| {
let db = connection.lock().unwrap();
let e = index_episode(&db, &x);
.collect();
let conn = db.lock().unwrap();
let e = conn.transaction::<(), Error, _>(|| {
episodes.iter().for_each(|x| {
let e = index_episode(&conn, &x);
if let Err(err) = e {
error!("Failed to index episode: {:?}.", x);
error!("Error msg: {}", err);
};
});
Ok(())
});
drop(conn);
if let Err(err) = e {
error!("Episodes Transcaction Failed.");
error!("Error msg: {}", err);
};
}
// Maybe this can be refactored into an Iterator for lazy evaluation.
pub fn fetch_feeds(connection: &Arc<Mutex<SqliteConnection>>, force: bool) -> Result<Vec<Feed>> {
let tempdb = connection.lock().unwrap();
let mut feeds = dbqueries::get_sources(&tempdb)?;
drop(tempdb);
pub fn fetch_feeds(db: &Arc<Mutex<SqliteConnection>>, force: bool) -> Result<Vec<Feed>> {
let mut feeds = {
let conn = db.lock().unwrap();
dbqueries::get_sources(&conn)?
};
let results: Vec<Feed> = feeds
.par_iter_mut()
.filter_map(|x| {
let db = connection.lock().unwrap();
let l = refresh_source(&db, x, force);
let l = refresh_source(db, x, force);
if l.is_ok() {
l.ok()
} else {
@@ -172,7 +181,7 @@ pub fn fetch_feeds(connection: &Arc<Mutex<SqliteConnection>>, force: bool) -> Re
}
pub fn refresh_source(
connection: &SqliteConnection,
db: &Arc<Mutex<SqliteConnection>>,
feed: &mut Source,
force: bool,
) -> Result<Feed> {
@@ -212,7 +221,7 @@ pub fn refresh_source(
// _ => (),
// };
feed.update_etag(connection, &req)?;
feed.update_etag(db, &req)?;
Ok(Feed(req, feed.clone()))
}
+9 -2
View File
@@ -3,6 +3,8 @@ use diesel::SaveChangesDsl;
use SqliteConnection;
use reqwest::header::{ETag, LastModified};
use std::sync::{Arc, Mutex};
use schema::{episode, podcast, source};
use errors::*;
@@ -178,7 +180,11 @@ impl<'a> Source {
/// Extract Etag and LastModifier from req, and update self and the
/// corresponding db row.
pub fn update_etag(&mut self, con: &SqliteConnection, req: &reqwest::Response) -> Result<()> {
pub fn update_etag(
&mut self,
db: &Arc<Mutex<SqliteConnection>>,
req: &reqwest::Response,
) -> Result<()> {
let headers = req.headers();
// let etag = headers.get_raw("ETag").unwrap();
@@ -191,7 +197,8 @@ impl<'a> Source {
{
self.http_etag = etag.map(|x| x.tag().to_string().to_owned());
self.last_modified = lmod.map(|x| format!("{}", x));
self.save_changes::<Source>(con)?;
let con = db.lock().unwrap();
self.save_changes::<Source>(&*con)?;
}
Ok(())