How To Ingest App Metrics from Slack into Elasticsearch
Recently I started helping Cardbucks, a very early-stage startup team. They’re running pretty bare-bones during their early stage market-fit experiments and haven’t setup any application monitoring or business intelligence solution for their users yet. However, they’ve been logging all user actions to a Slack room from Day One, which is awesome. So for a hack day, I built a bot to scrape the historical messages as well as ingest all new incoming metrics from Slack into Elasticsearch.
Ingest Real-Time Metrics
The first thing was to find an easy bot framework that both lets me receive new messages in (near) real-time. The Slack Team have generously provided the skeleton with python-rtmbot. This is a callback-based bot engine, so we need only write a simple plugin and configure it with our Slack token for ingesting metrics from the real-time message stream.
Below is a simple example of how we did this for the Cardbucks team. Obviously, since we’re parsing the textual messages into meaningful metrics, you’ll need to customize this to the metrics messages that your application is logging to Slack.
"""<br />
Listens to incoming Slack messages from cardbucks-bot and stores the metrics in Elasticsearch<br />
"""<br />
from elasticsearch import Elasticsearch<br />
from datetime import datetime<br />
import re</p>
<p># rtmbot variables<br />
crontable = []<br />
outputs = []</p>
<p>ELASTIC_HOST = '[your-hostname]'<br />
ELASTIC_INDEX = '[your-index]'<br />
ELASTIC_TYPE = '[your-type]'<br />
USER_ID = '[your Slack bot id]'<br />
CHANNEL = '[your Slack channel id]'</p>
<p>METRIC_PATTERNS = {<br />
'user_registration': re.compile('New User Registration')<br />
'user_disabled': re.compile('User Disabled(: )?(?P.*)'),<br />
'user_enabled': re.compile('User Enabled'),<br />
'payment_failed': re.compile('Payment Failed'),<br />
'fingerprint_blacklisted': re.compile('Fingerprint Blacklisted'),<br />
'ip_addr_blacklisted': re.compile('IP Address Blacklisted'),<br />
'purchase_confirmation': re.compile('User purchase\n(?P.*)', re.DOTALL)<br />
# and many more types of events captured...<br />
}</p>
<p>POST_PROCESSING = {<br />
'purchases': lambda purchases, fields: process_purchases(purchases, fields)<br />
# ... and more of these too<br />
}</p>
<p># global variables<br />
client = Elasticsearch(ELASTIC_HOST)<br />
epoch = datetime.utcfromtimestamp(0)</p>
<p>def process_purchases(purchases, fields):<br />
if purchases:<br />
fields["cards"] = []<br />
for per_merchant in filter(None, purchases.split('\n')):<br />
count, merchant = per_merchant.strip().split(None, 1)<br />
fields["cards"].append({"merchant": merchant, "value": int(count)})<br />
return fields</p>
<p>def to_metric(data):<br />
for event, pattern in METRIC_PATTERNS.iteritems():<br />
m = pattern.match(data['text'])<br />
if m:<br />
doc = {<br />
'event': event,<br />
'source': 'slack',<br />
'timestamp': datetime.utcfromtimestamp(float(data['ts']))<br />
}<br />
fields = m.groupdict()<br />
for field, handler in POST_PROCESSING.iteritems():<br />
fields = handler(fields.pop(field, None), fields)<br />
return dict(doc.items() + fields.items())</p>
<p>def process_message(data):<br />
# Filter out only the metrics messages we care about<br />
if data['user'] == USER_ID and data['channel'] == CHANNEL:<br />
metric = to_metric(data)<br />
if not metric:<br />
print "ERROR parsing %s" % data<br />
else:<br />
client.index(ELASTIC_INDEX, ELASTIC_TYPE, metric)<br />
Backfill Historical Metrics
After we have new metrics coming into the system, we also need to backfill all of the existing metrics from Day One into Elasticsearch. This will help us get a much better idea of growth rate and visualize other trends in our data.
Reading through the Slack API, we see the `channels.history` api call. This looks like exactly what we need. Again, the Slack Team provides for us. The standard python-slackclient makes this such API calls easy to access.
#!/usr/bin/env python<br />
from slackclient import SlackClient<br />
from elasticsearch import Elasticsearch<br />
from plugins.elastic.elastic import CHANNEL, USER_ID, process_message, to_metric<br />
from datetime import datetime<br />
import yaml<br />
import json</p>
<p>ELASTIC_HOST = '[your-hostname]'<br />
SIZE = 1000<br />
OLDEST = 0<br />
LATEST = datetime.now()</p>
<p>config = yaml.load(file('rtmbot.conf', 'r'))<br />
sc = SlackClient(config["SLACK_TOKEN"])<br />
client = Elasticsearch(ELASTIC_HOST)</p>
<p>response = json.loads(sc.api_call('channels.history', channel=CHANNEL, count=SIZE, oldest=OLDEST, latest=LATEST))<br />
for message in response['messages']:<br />
if 'channel' not in message:<br />
message['channel'] = CHANNEL<br />
if 'user' not in message:<br />
message['user'] = USER_ID<br />
process_message(message)</p>
<p>while response['has_more']:<br />
oldest_ts = response['messages'][-1]['ts']<br />
response = json.loads(sc.api_call('channels.history', channel=CHANNEL, count=SIZE, oldest=OLDEST, latest=oldest_ts))<br />
for message in response['messages']:<br />
if 'channel' not in message:<br />
message['channel'] = CHANNEL<br />
if 'user' not in message:<br />
message['user'] = USER_ID<br />
process_message(message)</p>
<p>if 'is_limited' in response and response['is_limited']:<br />
print "We have more messages that we can't access. :'("<br />
Hopefully, you’re either early enough in the process, have few enough metrics, or pay for Slack so that you don’t run into their 10,000 message history limit.
Now that you have all of your data in Elasticsearch, go forth and visualize with Kibana!
1 comment
Imported from the previous site.