file_fetcher.py 2.0 KB

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