diff --git a/.gitignore b/.gitignore index 9e41ecc..8057737 100644 --- a/.gitignore +++ b/.gitignore @@ -1,3 +1,5 @@ log nohup.out *.swp +*.log +*.out diff --git a/consumer.py b/consumer.py new file mode 100755 index 0000000..2f700d8 --- /dev/null +++ b/consumer.py @@ -0,0 +1,50 @@ +import redis +import urllib2 +import logging +import MySQLdb +import time + +r = redis.Redis(host="192.168.80.55",port=6379,db=0) + +logger = logging.getLogger() +hdlr = logging.FileHandler("consumer.log") +formatter = logging.Formatter('%(asctime)s %(levelname)s %(message)s') +hdlr.setFormatter(formatter) +logger.addHandler(hdlr) +logger.setLevel(logging.NOTSET) +send_headers = {"Content-Type":"application/json","Authorization":""} +conn = MySQLdb.connect(host="192.168.80.104",user="influx",passwd="influx1234",db="pages",charset='utf8' ) +cursor = conn.cursor() +sql_add_repo = "insert into github_html_detail(url,detail,time) values(%s,%s,%s)" +sql_update_state = "update repo_urls set state=1 where id=%s" +def start(): + logger.info("consumer begains to work") + task = r.blpop("repo_urls")[1] + while True: + logger.info("task : %s",task) + id_to_set_state = task[0:task.index("-")] + url_download = task[task.index("-")+1:] + logger.info(">>id:%s - url:%s"%(id_to_set_state,url_download)) + + token = r.blpop("tokens")[1] + r.rpush("tokens",token) + logger.info(">>token:%s will be used"%token) + + send_headers["Authorization"] = "token %s"%token + req = urllib2.Request(url_download,headers = send_headers) + try: + response = urllib2.urlopen(req,timeout=20) + repo = response.read() + logger.info(">>done downloading") + cursor.execute(sql_add_repo,(url_download,repo,time.strftime('%Y-%m-%d %H:%M:%S'))) + cursor.execute(sql_update_state,(id_to_set_state,)) + conn.commit() + logger.info(">>done storing") + except Exception,e: + logger.error(e.message) + continue + task = r.blpop("repo_urls")[1] + +if __name__ == "__main__": + start() + diff --git a/list_repos.py b/list_repos.py index 04c12da..f55d0aa 100644 --- a/list_repos.py +++ b/list_repos.py @@ -18,16 +18,15 @@ sql_insert = "insert into repo_urls(link,url,full_name,state) values(%s,%s,%s,%s sql_error = "insert into errors(url,message,time) values(%s,%s,%s)" send_headers = {"Content-Type":"application/json","Authorization":"token e41a69a15db215bd7959320360e77a988c4e19a3"} -url = "https://api.github.com/repositories?since=330001" +url = "https://api.github.com/repositories?since=444679" while True: try: req = urllib2.Request(url,headers = send_headers) logger.info("process link :%s"%url) - result = urllib2.urlopen(req) - logger.info("downloading is done") + result = urllib2.urlopen(req,timeout=10) prjs = json.loads(result.read()) - logger.info("done parsing json") + logger.info("done downloading and parsing json") for prj in prjs: try: cursor.execute(sql_insert,(url,prj.get("url"),prj.get("full_name"),0)) diff --git a/producer.py b/producer.py old mode 100644 new mode 100755 index a5318d8..47aa09e --- a/producer.py +++ b/producer.py @@ -1,2 +1,48 @@ import redis -print redis.__doc__ +import logging +import MySQLdb +import time + +r = redis.Redis(host="192.168.80.55",port=6379,db=0) + +logger = logging.getLogger() +hdlr = logging.FileHandler("produce.log") +formatter = logging.Formatter('%(asctime)s %(levelname)s %(message)s') +hdlr.setFormatter(formatter) +logger.addHandler(hdlr) +logger.setLevel(logging.NOTSET) + +conn = MySQLdb.connect(host="192.168.80.104",user="influx",passwd="influx1234",db="pages",charset='utf8' ) +cursor = conn.cursor() +sql_select_pointer = "select pointer from pointers where table_name='repo_urls'" +sql_update_pointer = "update pointers set pointer=%s where table_name='repo_urls'" +sql_add_repo = "insert into github_html_detail(url,detail,time) values(%s,%s,%s)" +sql_update_pushtime = "update repo_urls set push_time=%s where id=%s" +sql_select_url = "select id,url from repo_urls where id>%s limit %s" +patch_size = 1000 #how many urls push to redis every time +def start(): + logger.info("producer begains to work") + while True: + cursor.execute(sql_select_pointer) + last_id = cursor.fetchone()[0] + logger.info(">>id after last process: %d",last_id) + + num_redis = r.llen("repo_urls") + logger.info(">>%d urls in redis"%num_redis) + if num_redis < 1000: + logger.info(">>add urls to redis") + cursor.execute(sql_select_url,(last_id,patch_size)) + urls_added = cursor.fetchall() + for url in urls_added: + r.rpush("repo_urls","%d-%s"%(url[0],url[1])) + cursor.execute(sql_update_pushtime,(time.time(),url[0])) + cursor.execute(sql_update_pointer,(last_id + patch_size,)) + conn.commit() + else: + logger.info(">>enough urls") + + time.sleep(10) + +if __name__ == "__main__": + start() + diff --git a/token.py b/token.py new file mode 100644 index 0000000..b7a2f7e --- /dev/null +++ b/token.py @@ -0,0 +1,18 @@ +import redis +import logging + + +logger = logging.getLogger() +hdlr = logging.FileHandler("token.log") +formatter = logging.Formatter('%(asctime)s %(levelname)s %(message)s') +hdlr.setFormatter(formatter) +logger.addHandler(hdlr) +logger.setLevel(logging.NOTSET) + +r = redis.Redis(host="192.168.80.55",port=6379,db=0) + +with open("tokens","r") as file: + for line in file.readlines(): + logger.info("add token: %s",line.strip()) + r.rpush("tokens",line.strip()) + diff --git a/tokens b/tokens new file mode 100644 index 0000000..31fef38 --- /dev/null +++ b/tokens @@ -0,0 +1,18 @@ +c9b850151b3aeb6d750619653bdd25ad27e79e76 +a36b692adab5f4bfa8d03fba34669b78fe52d134 +e41a69a15db215bd7959320360e77a988c4e19a3 +d1e7a9d2eb0f54ec0e96ca06a64d68897d122513 +cd1c3cda5c269c30ca3527968ea1a4e41038d8f1 +25826285cf58b8bea2712298f2bc1dc598c8200b +277152ee058361c684ee253b881141bba722dfc5 +7933edb46487203f2752a93c702fd630cd552bae +9377749a42d8492428c52329624a8388b620e926 +d510a45755084273b25a12703eb5c39245cf4db3 +c023106576987273c64d45c8252cbc1785460917 +1473e930e6a7674673b0e864711d12b2c38db5fd +2fcd8e5c4fcf05eea2fecd5bd617d4910c37881f +0e862a968c0f7ba84acc9b909c3fccdf3e9cd15a +877738b0ede13b627605e301dd4f00725697ca0d +8368357f10e6318309b7e278b900e375f73421bd +6ac08b18aa04a36b602957be808a813900487173 +e94b02919d147be8df7dfc5ede4de92c3a75b1b4