file_fetcher.py 2.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263
  1. import json
  2. import socket
  3. import sys
  4. import time
  5. from threading import Thread
  6. from bookworm.constants import REDIS_FETCH_FILE, REDIS_STATE_KEY, REDIS_STEP_KEY
  7. from bookworm.logger import log, setup_logger
  8. from bookworm import s3
  9. import redis
  10. def netcat(filename, ip, port, size, job_key, s3client, redis, meta):
  11. log.info('netcat: %s %d %d %s', ip, port, size, filename)
  12. s = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
  13. log.info("Fetching %s", filename)
  14. s.connect((ip, port))
  15. log.info("Receiving file %s", filename)
  16. redis.hset(job_key, REDIS_STEP_KEY, 'DOWNLOADING')
  17. buff = b''
  18. count = 0
  19. last_perc = 0
  20. while True:
  21. data = s.recv(16384)
  22. if len(data) == 0:
  23. log.info("No data received - finished")
  24. break
  25. count += len(data)
  26. buff += data
  27. perc = int(100 * count / size)
  28. if perc % 5 == 0 and perc != last_perc:
  29. log.info("Download percentage: %d", perc)
  30. redis.hset(job_key, REDIS_STATE_KEY, str(perc))
  31. last_perc = perc
  32. if count >= size:
  33. break
  34. log.info("Download complete")
  35. s.close()
  36. log.info("Putting file in s3")
  37. redis.hset(job_key, REDIS_STATE_KEY, '100')
  38. s3client.put_object(Body=buff, Bucket=meta['raw_file_bucket'], Key=filename)
  39. log.info("File %s in s3 with key %s", filename, job_key)
  40. redis.rpush(meta['unpack_file_queue'], json.dumps({'job_key': job_key, 's3key': filename, 'meta': meta}))
  41. def main():
  42. r = redis.StrictRedis(host='localhost', port=6379)
  43. setup_logger()
  44. while True:
  45. log.info('Waiting for message on %s', REDIS_FETCH_FILE)
  46. topic, message = r.blpop(REDIS_FETCH_FILE)
  47. log.info('got message: %s', message)
  48. params = json.loads(message.decode('utf-8'))
  49. log.info('params for netcat: %s', params)
  50. params['s3client'] = s3.client()
  51. params['redis'] = r
  52. t = Thread(target=netcat, kwargs=params)
  53. t.daemon = True
  54. t.start()
  55. main()