(self, query, num_streams)
| 44 | calling start(). |
| 45 | """ |
| 46 | def __init__(self, query, num_streams): |
| 47 | self.query = query |
| 48 | self.num_streams = num_streams |
| 49 | self.stop_ev = Event() |
| 50 | self.output_q = Queue() |
| 51 | self.threads = [] |
| 52 | self.query_rate = 0 |
| 53 | self.query_rate_thread = Thread(target=self.compute_query_rate, |
| 54 | args=(self.output_q, self.stop_ev)) |
| 55 | |
| 56 | def execute(self, query): |
| 57 | """Executes a query on the coordinator of the local minicluster.""" |