parser.py 3.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131
  1. #!/usr/bin/python3
  2. import threading
  3. from queue import Queue
  4. from db import db
  5. from source import Source
  6. from linklist import LinkList
  7. from news import News
  8. from paid_parser import PaidParser
  9. from time import strftime,time
  10. from multiprocessing import Pool
  11. import sys
  12. from analytics import post_data
  13. def parseSource(s):
  14. init = time()
  15. q = Queue()
  16. sel = None
  17. rss = False
  18. if "selector" in s:
  19. sel = s["selector"]
  20. if "rss" in s:
  21. rss = bool(s["rss"])
  22. l = LinkList(s["link"], sel, rss=rss)
  23. d = db()
  24. kw = d.get_keywords(s["category"])
  25. new_articles = 0
  26. if l is not None:
  27. for a in l.links:
  28. if not d.article_exists(a.lower(), s["_id"]):
  29. q.put(a)
  30. # print("[%s][Source] %s: %d/%d new/total articles" %
  31. # (strftime("%H:%M:%S"), s["name"], q.qsize(), len(l.links)))
  32. new_articles = q.qsize()
  33. if not q.empty():
  34. num_worker_threads = min(5, q.qsize())
  35. threads = []
  36. for _ in range(num_worker_threads):
  37. t = threading.Thread(target=news_worker, args=[s, kw, d, q])
  38. # source, keywords, db, queue
  39. t.daemon = True
  40. t.start()
  41. threads.append(t)
  42. q.join() # block until all tasks are done
  43. for _ in range(num_worker_threads):
  44. # put num_worker empty jobs so
  45. # each worker can work once
  46. q.put(None)
  47. for t in threads:
  48. t.join(10)
  49. ex_time = time()-init
  50. print("[%s][Source] %s => Finished parsing, %d/%d new articles, in %0.3fsec" %
  51. (strftime("%H:%M:%S"),
  52. s["name"], new_articles, len(l.links), ex_time))
  53. sys.stdout.flush()
  54. return (ex_time, len(l.links), s["name"], new_articles)
  55. def news_worker(source, kw, d, q):
  56. source_id = source["_id"]
  57. source_name = source["name"]
  58. source_category = source["category"]
  59. paid = False
  60. if "paid" in source and source["paid"]:
  61. paid = True
  62. # tname = threading.current_thread().name
  63. while True:
  64. try:
  65. item = q.get(timeout=5)
  66. except:
  67. # print("Failed to get item")
  68. q.task_done()
  69. break
  70. if item is None:
  71. q.task_done()
  72. break
  73. # print("[Thread %s] %s" % (tname,item))
  74. try:
  75. html = None
  76. if paid:
  77. p = PaidParser(item)
  78. html = p.html
  79. n = News(item, source_id, source_name,
  80. source_category, kw, html=html)
  81. # print("[Thread %s] Finished" % tname)
  82. d.insert_article(n.get())
  83. except Exception as e:
  84. print("#### EXCEPTION ########")
  85. print(e)
  86. print(item)
  87. print("#### END EXCEPTION ####")
  88. finally:
  89. q.task_done()
  90. post_data('event', {}, '"start"')
  91. start = time()
  92. d = db()
  93. sources = d.sources()
  94. post_data('value', {'type': 'error'}, d.count_error())
  95. post_data('event', {}, '"purge_error_finish"')
  96. d.purge_error()
  97. post_data('event', {}, '"purge_error_finish"')
  98. pool = Pool(processes=4)
  99. ret = pool.map(parseSource, sources)
  100. new = sum([r[3] for r in ret])
  101. tot = sum([r[1] for r in ret])
  102. post_data('value', {'type': 'new'}, new)
  103. post_data('value', {'type': 'total'}, tot)
  104. # n = sum(ret)
  105. # sret = sorted(ret, key=lambda tup: tup[0])
  106. # print(sret)
  107. print("[%s] %d/%d new articles, total time %d sec" %
  108. (strftime("%H:%M:%S"), new, tot, time()-start))
  109. post_data('time', {'type': 'total'}, time()-start)
  110. print("#" * 75)
  111. post_data('event', {}, '"end"')