Repository navigation
Job only running one time #4
Description
Activity
UPDATES
I'm running this now with debugging. I started two workers, each running a seperate job. Now they both run once and then the second worker goes into a really weird non stop loop. Note: the loop is seen in the second snippet here below (redis debugger)
Hworker debugger:
*Main> ("notifier",("WORKER RUNNING","[\"450733c6-3c6f-451c-9a25-9a5570b19f85\",{\"tag\":\"ServiceAction\",\"contents\":[\"test\",\"testing\",1]}]")) "service:test|action:testing|id:1" ("notifier",("JOB COMPLETE","[\"450733c6-3c6f-451c-9a25-9a5570b19f85\",{\"tag\":\"ServiceAction\",\"contents\":[\"test\",\"testing\",1]}]")) ("notifier",("MONITOR RV",1)) ("notifier",("MONITOR REQUEUED","[\"df925ed6-4c86-4989-b967-99230d13715a\",{\"tag\":\"ServiceAction\",\"contents\":[\"test\",\"testing\",1]}]")) ("notifier",("WORKER RUNNING","[\"df925ed6-4c86-4989-b967-99230d13715a\",{\"tag\":\"ServiceAction\",\"contents\":[\"test\",\"testing\",1]}]")) ("search",("WORKER RUNNING","[\"38410142-6f83-4348-8f95-d2ab1dd19570\",{\"tag\":\"SearchIndex\",\"contents\":[\"test\",2]}]")) ("search",("JOB COMPLETE","[\"38410142-6f83-4348-8f95-d2ab1dd19570\",{\"tag\":\"SearchIndex\",\"contents\":[\"test\",2]}]")) ("notifier",("MONITOR RV",1)) ("notifier",("MONITOR REQUEUED","[\"df925ed6-4c86-4989-b967-99230d13715a\",{\"tag\":\"ServiceAction\",\"contents\":[\"test\",\"testing\",1]}]"))Redis debug
1513173783.022352 [0 172.17.0.1:46054] "EVAL" "local job = redis.call('rpop',KEYS[1])\nif job ~= nil then\n redis.call('hset', KEYS[2], job, ARGV[1])\n return job\nelse\n return nil\nend" "2" "hworker-jobs-search" "hworker-progress-search" "\"2017-12-13T14:03:03.021753Z\"" 1513173783.022405 [0 lua] "rpop" "hworker-jobs-search" 1513173783.034465 [0 172.17.0.1:46054] "EVAL" "local job = redis.call('rpop',KEYS[1])\nif job ~= nil then\n redis.call('hset', KEYS[2], job, ARGV[1])\n return job\nelse\n return nil\nend" "2" "hworker-jobs-search" "hworker-progress-search" "\"2017-12-13T14:03:03.03398Z\"" 1513173783.034511 [0 lua] "rpop" "hworker-jobs-search" 1513173783.046367 [0 172.17.0.1:46054] "EVAL" "local job = redis.call('rpop',KEYS[1])\nif job ~= nil then\n redis.call('hset', KEYS[2], job, ARGV[1])\n return job\nelse\n return nil\nend" "2" "hworker-jobs-search" "hworker-progress-search" "\"2017-12-13T14:03:03.045846Z\"" 1513173783.046417 [0 lua] "rpop" "hworker-jobs-search" 1513173783.057326 [0 172.17.0.1:46054] "EVAL" "local job = redis.call('rpop',KEYS[1])\nif job ~= nil then\n redis.call('hset', KEYS[2], job, ARGV[1])\n return job\nelse\n return nil\nend" "2" "hworker-jobs-search" "hworker-progress-search" "\"2017-12-13T14:03:03.056814Z\"" 1513173783.057373 [0 lua] "rpop" "hworker-jobs-search" 1513173783.070731 [0 172.17.0.1:46054] "EVAL" "local job = redis.call('rpop',KEYS[1])\nif job ~= nil then\n redis.call('hset', KEYS[2], job, ARGV[1])\n return job\nelse\n return nil\nend" "2" "hworker-jobs-search" "hworker-progress-search" "\"2017-12-13T14:03:03.069905Z\"" 1513173783.070801 [0 lua] "rpop" "hworker-jobs-search"- Please update if you are still running into issues, but in answer to your last question, you can have as many queues as you want on the same redis server (they each need a different name — i.e., you named yours “notifier”).…On Dec 13, 2017, at 9:39 AM, Donna ***@***.***> wrote: Does "A reliable at-least-once job queue built on Redis." mean that it can only run one queue per redis database? — You are receiving this because you are subscribed to this thread. Reply to this email directly, view it on GitHub, or mute the thread.
Hi @dbp thank you for answering. I'm trying to see what might be wrong here. I did name my queues differently. Now I have this
let config_1 = (defaultHworkerConfig "notifier" (State state)) { hwconfigDebug = True} let config_2 = (defaultHworkerConfig "search" (State state)) { hwconfigDebug = True} -- Start the worker nworker <- createWith config_1 sworker <- createWith config_2 forkIO (monitor nworker) forkIO (worker nworker) forkIO (worker sworker) forkIO (monitor sworker) forkIO (forever $ queue nworker (ServiceAction "test" "testing" 1) >> threadDelay 5000000) forkIO (forever $ queue sworker (SearchIndex "test" 2) >> threadDelay 5000000)On my first run after changing this, the queue named search worked but on a second run it didn't work anymore. I have the debugger running and occasionally the monitor says "MONITOR REQUEUED" but then nothing. I'm writing the results of both jobs to a file so I should see a new line of text when it runs.
Could this be because of the threads I create? I'm very new to threading in Haskell
The redis debugger shows that it's definitely pushing the tasks into the queue but does not seem to process them.
1513180391.971306 [0 172.17.0.1:46532] "HKEYS" "hworker-progress-search" 1513180393.126282 [0 172.17.0.1:46532] "LPUSH" "hworker-jobs-search" "[\"fc3f116c-4ba3-46b6-b233-7f6fca9512d3\",{\"tag\":\"SearchIndex\",\"contents\":[\"test\",2]}]" 1513180398.130543 [0 172.17.0.1:46532] "LPUSH" "hworker-jobs-search" "[\"8061a508-6a84-4fe1-b810-dbbd51d2d3f4\",{\"tag\":\"SearchIndex\",\"contents\":[\"test\",2]}]" 1513180403.136544 [0 172.17.0.1:46532] "LPUSH" "hworker-jobs-search" "[\"997a9d6d-fa3a-427b-9d18-b94eda986cfc\",{\"tag\":\"SearchIndex\",\"contents\":[\"test\",2]}]" 1513180403.974982 [0 172.17.0.1:46532] "HKEYS" "hworker-progress-search" 1513180408.137567 [0 172.17.0.1:46532] "LPUSH" "hworker-jobs-search" "[\"d8af7c59-9357-4d0f-b80a-4694e6e248c3\",{\"tag\":\"SearchIndex\",\"contents\":[\"test\",2]}]" 1513180413.141071 [0 172.17.0.1:46532] "LPUSH" "hworker-jobs-search" "[\"de15be50-058f-432e-97cd-d7434f2197f8\",{\"tag\":\"SearchIndex\",\"contents\":[\"test\",2]}]" 1513180415.978862 [0 172.17.0.1:46532] "HKEYS" "hworker-progress-search" 1513180418.145018 [0 172.17.0.1:46532] "LPUSH" "hworker-jobs-search" "[\"f9dd7014-8c50-419b-83d1-9ec5f9359e17\",{\"tag\":\"SearchIndex\",\"contents\":[\"test\",2]}]" 1513180423.150860 [0 172.17.0.1:46532] "LPUSH" "hworker-jobs-search" "[\"be5233c4-b648-48cd-85bf-7b18ca1699d9\",{\"tag\":\"SearchIndex\",\"contents\":[\"test\",2]}]" 1513180427.982957 [0 172.17.0.1:46532] "HKEYS" "hworker-progress-search" 1513180428.155229 [0 172.17.0.1:46532] "LPUSH" "hworker-jobs-search" "[\"181b57ce-eedd-4e94-9ea7-f0c1f8ab002b\",{\"tag\":\"SearchIndex\",\"contents\":[\"test\",2]}]"This is what makes me think it has something to do with threading. The debugger is showing garbled output
*Main> ("search",("W(O"RnKoEtRi fRiUeNrN"I,N(G""W,O"R[K\E"Re dR8UfN0N9IdN4G-"3,2"a[3\-"41ba5c6f-3a761331--6382f0a0-f4439188b-b8727b\8"-,9{b\5"6t8a3ga\c"3:f\b"bS\e"r,v{i\c"etAacgt\i"o:n\\""S,e\r"vciocnetAecnttiso\n"\:"[,\\""tceosntt\e"n,t\s"\t"e:s[t\i"ntge\s"t,\1"],}\]""t)e)s ting\",1]}(]""s)e)a rch",("JOB COMPLETE","[\"ed8f09d4-32a3-4b56-a631-3200f498bb77\",{\"tag\":\"ServiceAction\",\"contents\":[\"test\",\"testing\",1]}]")) ("search",("WORKER RUNNING","[\"cab85ce4-ba73-400a-8f86-22144ac0fca7\",{\"tag\":\"ServiceAction\",\"contents\":[\"test\",\"testing\",1]}]"))So based on the program fragment you gave above, one problem might be that every thread you forked off with
forkIO, which means that the main thread will eventually terminate, which will then kill the program. e.g., the following program will print (on my machine) about 1200 times and then quit:import Control.Concurrent (forkIO) import Control.Monad (forever) main :: IO () main = do forkIO $ forever (putStrLn "Echo") return ()Usually the way around this is to leave one of them (usually the most important one, in this case perhaps the monitor?) in the main thread (i.e., just leave off the
forkIOcall). If that's not it, is there a way you can post a little more code?I tried this but I get the same results. Here is my entire code. Don't mind the Websocket stuff it's not really being used at the moment
{-# LANGUAGE MultiParamTypeClasses #-} {-# LANGUAGE DeriveGeneric #-} {-# LANGUAGE OverloadedStrings #-} module Main where import Data.List (intercalate) import GHC.Generics (Generic) import Data.Aeson (FromJSON, ToJSON) import Control.Concurrent.MVar (MVar, newMVar, putMVar, takeMVar) import Control.Concurrent (forkIO, threadDelay) import Control.Monad (forever, forM_) import System.Hworker import qualified Data.Text as T import qualified Network.WebSockets as WS type Client = (T.Text, WS.Connection) type ServerState = [Client] data ServiceJob = ServiceAction String String Int | SearchIndex String Int deriving (Generic, Show) data State = State (MVar ServerState) instance ToJSON ServiceJob instance FromJSON ServiceJob instance Job State ServiceJob where job (State mvar) (ServiceAction service action id) = do clients <- takeMVar mvar let message = "service:" ++ (service) ++ "|action:" ++ (action) ++ "|id:" ++ show id appendFile "/Users/Donna/Code/K0TT/repos/cherry-service-worker-test/test.txt" (message ++ "\n") >> return Success job (State mvar) (SearchIndex service id) = do appendFile "/Users/Donna/Code/K0TT/repos/cherry-service-worker-test/test.txt" "search\n" >> return Success broadcast :: T.Text -> ServerState -> IO () broadcast message clients = do forM_ clients $ \(_, conn) -> WS.sendTextData conn message newServerState :: ServerState newServerState = [] main = do state <- newMVar newServerState let config_1 = (defaultHworkerConfig "notifier" (State state)) { hwconfigDebug = True} let config_2 = (defaultHworkerConfig "search" (State state)) { hwconfigDebug = True} nworker <- createWith config_1 sworker <- createWith config_2 -- Start the notifier worker forkIO (worker nworker) -- Start the search worker forkIO (worker sworker) forkIO (monitor sworker) -- Add new ServiceAction every 5 secs forkIO (forever $ queue nworker (ServiceAction "test" "testing" 1) >> threadDelay 5000000) -- Add new SearchIndex every 10 secs forkIO (forever $ queue sworker (SearchIndex "test" 2) >> threadDelay 10000000) -- Monitor the notifier worker monitor nworkerSo there are two things I notice:
-
Your ServiceAction takes the MVar but never refills it, so this will lock the thread the second time it runs (MVars are either empty or full, and different functions have different blocking behavior --
takeMVarwaits until it gets filled if it is empty, but since that'll never happen here, it'll wait forever). -
It shouldn't matter here, but you don't actually need two separate queues (if you did that intentionally, ignore me!). A given queue can store a single data type, but it can have multiple variants (in this case, ServiceActions and SearchIndexes). What you are doing instead is creating two different queues, both of which can store either type of job, and then only storing one type in each queue. If you wanted to separate them completely (because in the future you might want the ServiceActions to run on one server and the SearchIndexes on another), it might be better to create different data types for them, so that way you couldn't accidentally queue the wrong type of job into a queue.
I ran your program locally (with minor changes for the path of where it writes the log messages), and I see the behavior I expect based on the above -- the ServiceAction task runs exactly once (as it then gets blocked), and the SearchIndex runs every ten seconds. You are using redis 2.8 right? We really need to support newer versions, but they changed some stuff in EVAL which we rely heavily on!
-
Ahh of course, that was it. I was not paying attention to the MVar because I wasn't using it yet. Thank you so much and sorry to waste your time with this. And you are correct, I don't need two queues but I thought that maybe that was the reason for the delay of the queues. Brilliant, now everything works. Thanks again 👍
Glad it's working, and no worries -- happy I was able to help!
Reacted by Donna- added a commit that references this issue
on Dec 30, 2024
I'm wondering why my queue is not getting processed. When I run this code a single job is processed but then if I add debugging I see that the queue is just filling up. My code
This should push a new job into the queue every 5 seconds, which it does. But it's not processing it. I checked the failed jobs queue and it's completely empty so they are not failing.
This is the Job instance