49 lines
1.6 KiB
Python
Executable File
49 lines
1.6 KiB
Python
Executable File
import redis
|
|
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 + len(urls_added),))
|
|
conn.commit()
|
|
else:
|
|
logger.info(">>enough urls")
|
|
|
|
time.sleep(10)
|
|
|
|
if __name__ == "__main__":
|
|
start()
|
|
|