producer and consumer is done
This commit is contained in:
parent
7b01001b9d
commit
5d531622f5
|
|
@ -1,3 +1,5 @@
|
|||
log
|
||||
nohup.out
|
||||
*.swp
|
||||
*.log
|
||||
*.out
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
||||
|
|
@ -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))
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
||||
|
|
|
|||
|
|
@ -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())
|
||||
|
||||
|
|
@ -0,0 +1,18 @@
|
|||
c9b850151b3aeb6d750619653bdd25ad27e79e76
|
||||
a36b692adab5f4bfa8d03fba34669b78fe52d134
|
||||
e41a69a15db215bd7959320360e77a988c4e19a3
|
||||
d1e7a9d2eb0f54ec0e96ca06a64d68897d122513
|
||||
cd1c3cda5c269c30ca3527968ea1a4e41038d8f1
|
||||
25826285cf58b8bea2712298f2bc1dc598c8200b
|
||||
277152ee058361c684ee253b881141bba722dfc5
|
||||
7933edb46487203f2752a93c702fd630cd552bae
|
||||
9377749a42d8492428c52329624a8388b620e926
|
||||
d510a45755084273b25a12703eb5c39245cf4db3
|
||||
c023106576987273c64d45c8252cbc1785460917
|
||||
1473e930e6a7674673b0e864711d12b2c38db5fd
|
||||
2fcd8e5c4fcf05eea2fecd5bd617d4910c37881f
|
||||
0e862a968c0f7ba84acc9b909c3fccdf3e9cd15a
|
||||
877738b0ede13b627605e301dd4f00725697ca0d
|
||||
8368357f10e6318309b7e278b900e375f73421bd
|
||||
6ac08b18aa04a36b602957be808a813900487173
|
||||
e94b02919d147be8df7dfc5ede4de92c3a75b1b4
|
||||
Loading…
Reference in New Issue