浏览代码

add tasks/workers

David 10 年之前
父节点
当前提交
04a0f524ba
共有 4 个文件被更改,包括 74 次插入1 次删除
  1. 3 0
      .gitignore
  2. 24 0
      back/tasks.py
  3. 30 1
      back/web.py
  4. 17 0
      back/worker.py

+ 3 - 0
.gitignore

@@ -0,0 +1,3 @@
+*pyc
+*/__pycache__/*
+*.swp

+ 24 - 0
back/tasks.py

@@ -0,0 +1,24 @@
+import time
+from subprocess import Popen,PIPE
+from rq import get_current_job
+def long_job(arg):
+    for i in range(20):
+        time.sleep(1)
+        job = get_current_job()
+        job.meta['progress'] = i
+        job.save()
+    #push to db, return id if ok. return error if not ok
+    return "37"
+
+
+def run_process(process=None):
+    process=['ffmpeg', "-loglevel","error","-stats",'-i', '/home/david/movie-1467810207.flv', '-ss', '0','-t','145', '-c:v','libx264','-movflags','+faststart','-y','output.mp4']
+    proc=Popen(process, stderr=PIPE, universal_newlines=True)
+    while True:
+        line = proc.stderr.readline()
+        if line != '' and line !=b'':
+            job = get_current_job()
+            job.meta['progress'] = line.rstrip()
+            job.save()
+        else:
+            break

+ 30 - 1
back/web.py

@@ -4,12 +4,18 @@ import json
 import urllib
 from db import db
 from flask import Flask, abort, jsonify, request
+import time
+import tasks
+
+from rq import Queue
+from rq.job import Job
+from worker import conn
+q = Queue(connection=conn)
 
 PORT = 8981
 BASE_PATH="/back" #changes based on webserver
 app = Flask(__name__)
 
-
 @app.route(BASE_PATH + '/')
 def index():
     return "Hello world"
@@ -22,6 +28,28 @@ def raw_list():
         }
     return jsonify(ret)
 
+@app.route(BASE_PATH + '/job')
+def job():
+    job = q.enqueue_call(
+            func=tasks.long_job, args=("testarg",), result_ttl=5000
+    )
+    ret = {"id":job.get_id()}
+
+    return jsonify(ret)
+
+
+@app.route(BASE_PATH+"/results/<job_key>", methods=['GET'])
+def get_results(job_key):
+
+    print(job_key)
+    job = Job.fetch(job_key, connection=conn)
+
+    if job.is_finished:
+        return str(job.result), 200
+    else:
+        return "Nay! %s" % job.meta.get('progress') , 202
+
+
 @app.route(BASE_PATH + '/cut')
 def cut_list():
     ret=[]
@@ -33,5 +61,6 @@ def uploaded_list():
     return jsonify(ret)
 
 
+
 if __name__ == '__main__':
     app.run(host="0.0.0.0",port=PORT, debug=True)

+ 17 - 0
back/worker.py

@@ -0,0 +1,17 @@
+#!/usr/bin/env python3
+import os
+
+import redis
+from rq import Worker, Queue, Connection
+
+listen = ['default']
+
+redis_url = os.getenv('REDISTOGO_URL', 'redis://localhost:6379')
+
+conn = redis.from_url(redis_url)
+
+if __name__ == '__main__':
+    with Connection(conn):
+        worker = Worker(list(map(Queue, listen)))
+        worker.work()
+