-
Notifications
You must be signed in to change notification settings - Fork 15
/
Copy pathconsume_realtime_results.py
executable file
·166 lines (126 loc) · 5.55 KB
/
consume_realtime_results.py
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
#!/usr/bin/env python
# ----------------------------------------------------------------------
# Numenta Platform for Intelligent Computing (NuPIC)
# Copyright (C) 2015, Numenta, Inc. Unless you have purchased from
# Numenta, Inc. a separate commercial license for this software code, the
# following terms and conditions apply:
#
# This program is free software: you can redistribute it and/or modify
# it under the terms of the GNU General Public License version 3 as
# published by the Free Software Foundation.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE.
# See the GNU General Public License for more details.
#
# You should have received a copy of the GNU General Public License
# along with this program. If not, see http://www.gnu.org/licenses.
#
# http://numenta.org/licenses/
# ----------------------------------------------------------------------
"""Consume anomaly results in near realtime"""
import os
from nta.utils import amqp
from nta.utils.config import Config
from htmengine import htmengineerrno
from htmengine.runtime.anomaly_service import AnomalyService
appConfig = Config("application.conf", os.environ["APPLICATION_CONFIG_PATH"])
modelResultsExchange = appConfig.get("metric_streamer",
"results_exchange_name")
queueName = "skeleton_results"
def declareExchanges(amqpClient):
""" Declares model results and non-metric data exchanges
"""
amqpClient.declareExchange(exchange=modelResultsExchange,
exchangeType="fanout",
durable=True)
def declareQueueAndBindToExchanges(amqpClient):
""" Declares skeleton queue and binds to model results.
"""
result = amqpClient.declareQueue(queueName, durable=True)
amqpClient.bindQueue(exchange=modelResultsExchange,
queue=result.queue, routingKey="")
def configChannel(amqpClient):
amqpClient.requestQoS(prefetchCount=1)
def handleModelInferenceResults(body):
""" Model results batch handler.
:param body: Serialized message payload; the message is compliant with
htmengine/runtime/json_schema/model_inference_results_msg_schema.json.
:type body: str
"""
try:
batch = AnomalyService.deserializeModelResult(body)
except Exception:
print "Error deserializing model result"
raise
metricId = batch["metric"]["uid"]
metricName = batch["metric"]["name"]
print "Handling %d model result(s) for %s - %s" % (len(batch["results"]),
metricId,
metricName)
if not batch["results"]:
print "Empty results in model inference results batch; model=%s" % metricId
return
print metricId, batch["results"]
def handleModelCommandResult(body):
""" ModelCommandResult handler. Handles model creation/deletion events
:param body: Incoming message payload
:type body: str
"""
try:
modelCommandResult = AnomalyService.deserializeModelResult(body)
except Exception:
print "Error deserializing model command result"
raise
if modelCommandResult["status"] != htmengineerrno.SUCCESS:
return # Ignore...
if modelCommandResult["method"] == "defineModel":
print "Handling `defineModel` for %s" % modelCommandResult.get("modelId")
print modelCommandResult
elif modelCommandResult["method"] == "deleteModel":
print "Handling `deleteModel` for %s" % modelCommandResult.get("modelId")
print modelCommandResult
def messageHandler(message):
""" Inspect all inbound model results
We will key off of routing key to determine specific handler for inbound
message. If routing key is `None`, attempt to decode message using
`AnomalyService.deserializeModelResult()`.
:param amqp.messages.ConsumerMessage message: ``message.body`` is one of:
Serialized batch of model inference results generated in
``AnomalyService`` and must be deserialized using
``AnomalyService.deserializeModelResult()``. Per
htmengine/runtime/json_schema/model_inference_results_msg_schema.json
Serialized ``ModelCommandResult`` generated in ``AnomalyService``
per model_command_result_amqp_message.json and must be deserialized
using ``AnomalyService.deserializeModelResult()``
"""
if message.methodInfo.routingKey is None:
print "Unrecognized routing key."
else:
dataType = (message.properties.headers.get("dataType")
if message.properties.headers else None)
if not dataType:
handleModelInferenceResults(message.body)
elif dataType == "model-cmd-result":
handleModelCommandResult(message.body)
else:
print "Unexpected message header dataType=%s" % dataType
message.ack()
if __name__ == "__main__":
with amqp.synchronous_amqp_client.SynchronousAmqpClient(
amqp.connection.getRabbitmqConnectionParameters(),
channelConfigCb=configChannel) as amqpClient:
declareExchanges(amqpClient)
declareQueueAndBindToExchanges(amqpClient)
consumer = amqpClient.createConsumer(queueName)
# Start consuming messages
for evt in amqpClient.readEvents():
if isinstance(evt, amqp.messages.ConsumerMessage):
messageHandler(evt)
elif isinstance(evt, amqp.consumer.ConsumerCancellation):
# Bad news: this likely means that our queue was deleted externally
msg = "Consumer cancelled by broker: %r (%r)" % (evt, consumer)
raise Exception(msg)
else:
print "Unexpected amqp event=%r" % evt