rewrite load
This commit is contained in:
+23
-11
@@ -94,19 +94,11 @@ func LoadUnprocessed() {
|
|||||||
db := getDb()
|
db := getDb()
|
||||||
defer db.Close()
|
defer db.Close()
|
||||||
for {
|
for {
|
||||||
num := countTelegram()
|
telegrams := getTelegram(db)
|
||||||
log.Println("unprocessed : ", num)
|
if len(telegrams) < 1 {
|
||||||
if num < 1 {
|
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
log.Println("try to load unprocessed telegram ", num)
|
for _, t := range telegrams {
|
||||||
rows, err := db.Query(load)
|
|
||||||
checkDbErr(err, "error in query telegram")
|
|
||||||
defer rows.Close()
|
|
||||||
for rows.Next() {
|
|
||||||
var t telegram
|
|
||||||
err := rows.Scan(&t.id, &t.text)
|
|
||||||
checkDbErr(err, "error in scan telegram")
|
|
||||||
if WriteToClient(t.text) == nil {
|
if WriteToClient(t.text) == nil {
|
||||||
processed(db, t)
|
processed(db, t)
|
||||||
}
|
}
|
||||||
@@ -114,6 +106,26 @@ func LoadUnprocessed() {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func getTelegram(db *sql.DB) []telegram {
|
||||||
|
var telegrams []telegram
|
||||||
|
num := countTelegram()
|
||||||
|
log.Println("unprocessed : ", num)
|
||||||
|
if num < 1 {
|
||||||
|
return telegrams
|
||||||
|
}
|
||||||
|
log.Println("try to load unprocessed telegram ", num)
|
||||||
|
rows, err := db.Query(load)
|
||||||
|
checkDbErr(err, "error in query telegram")
|
||||||
|
defer rows.Close()
|
||||||
|
for rows.Next() {
|
||||||
|
var t telegram
|
||||||
|
err := rows.Scan(&t.id, &t.text)
|
||||||
|
checkDbErr(err, "error in scan telegram")
|
||||||
|
telegrams = append(telegrams, t)
|
||||||
|
}
|
||||||
|
return telegrams
|
||||||
|
}
|
||||||
|
|
||||||
func processed(db *sql.DB, t telegram) {
|
func processed(db *sql.DB, t telegram) {
|
||||||
getWriteLock()
|
getWriteLock()
|
||||||
tx, err := db.Begin()
|
tx, err := db.Begin()
|
||||||
|
|||||||
Reference in New Issue
Block a user