Dash에 연결하기
Dash에 연결하기 (Connect to Dash)
이 Apache Pinot 가이드에서는 Dash 웹 프레임워크를 사용해 데이터를 시각화하는 방법을 알려드려요.
이 가이드에서는 Plotly의 Dash 웹 프레임워크를 사용해 Apache Pinot의 데이터를 시각화하는 방법을 배워요. Dash는 ML 및 데이터 사이언스 웹 앱을 만들기 위한 가장 많이 다운로드되고 신뢰받는 Python 프레임워크예요.
Dash를 사용해 Wikimedia 속성에 가해지는 변경 사항을 시각화하는 실시간 대시보드를 만들 거예요.
출처: 문서
본문
시작 구성요소 (Startup components)
Zookeeper, Kafka 인스턴스와 함께 Pinot controller, broker, server를 띄우는 다음 Docker compose 파일을 사용할 거예요:
version: '3.7'
services:
zookeeper:
image: zookeeper:3.9.5
container_name: "zookeeper-wiki"
ports:
- "2181:2181"
environment:
ZOOKEEPER_CLIENT_PORT: 2181
ZOOKEEPER_TICK_TIME: 2000
kafka:
image: apache/kafka:4.0.0
restart: unless-stopped
container_name: "kafka-wiki"
ports:
- "9092:9092"
expose:
- "9093"
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9093,OUTSIDE://0.0.0.0:9092,CONTROLLER://0.0.0.0:29093
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka-wiki:9093,OUTSIDE://localhost:9092
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,OUTSIDE:PLAINTEXT
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka-wiki:29093
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
CLUSTER_ID: MkU3OEVBNTcwNTJENDM2Qk
pinot-controller:
image: apachepinot/pinot:1.5.1
command: "StartController -zkAddress zookeeper-wiki:2181 -dataDir /data"
container_name: "pinot-controller-wiki"
volumes:
- ./config:/config
- ./data:/data
restart: unless-stopped
ports:
- "9000:9000"
depends_on:
- zookeeper
pinot-broker:
image: apachepinot/pinot:1.5.1
command: "StartBroker -zkAddress zookeeper-wiki:2181"
restart: unless-stopped
container_name: "pinot-broker-wiki"
volumes:
- ./config:/config
ports:
- "8099:8099"
depends_on:
- pinot-controller
pinot-server:
image: apachepinot/pinot:1.5.1
command: "StartServer -zkAddress zookeeper-wiki:2181"
restart: unless-stopped
container_name: "pinot-server-wiki"
volumes:
- ./config:/config
depends_on:
- pinot-broker
docker-compose.yml
모든 구성요소를 실행하려면 다음 명령을 실행하세요:
docker-compose up
Wikimedia 최근 변경 스트림
Wikimedia는 다양한 Wikimedia 속성에 가해진 변경을 설명하는 구조화된 이벤트 데이터의 연속 스트림을 제공해요. 이벤트는 Server-Side Events (SSE) 프로토콜을 사용해 HTTP로 게시돼요.
엔드포인트는 다음에서 찾을 수 있어요: stream.wikimedia.org/v2/stream/recentchange
이 데이터를 소비하려면 SSE 클라이언트 라이브러리를 설치해야 해요:
pip install sseclient-py
다음으로 wiki.py라는 파일을 만들어 다음 내용을 담으세요:
import json
import pprint
import sseclient
import requests
def with_requests(url, headers):
"""Get a streaming response for the given event feed using requests."""
return requests.get(url, stream=True, headers=headers)
url = 'https://stream.wikimedia.org/v2/stream/recentchange'
headers = {'Accept': 'text/event-stream'}
response = with_requests(url, headers)
client = sseclient.SSEClient(response)
for event in client.events():
stream = json.loads(event.data)
pprint.pprint(stream)
wiki.py
강조된 섹션은 SSE 클라이언트 라이브러리를 사용해 최근 변경 피드에 연결하는 방법을 보여줘요.
이 스크립트를 아래와 같이 실행해보세요:
python wiki.py
(잘린) 다음 출력을 볼 거예요:
출력
{'$schema': '/mediawiki/recentchange/1.0.0',
'bot': False,
'comment': '[...]',
'id': 1923506287,
'meta': {'domain': 'commons.wikimedia.org', 'dt': '2022-05-12T09:57:00Z', 'id': '3800228e-43d8-440d-8034-c68977742653', 'offset': 3855767440, 'partition': 0, 'request_id': '930b17cc-f14a-4656-afa1-d15b79a8f666', 'stream': 'mediawiki.recentchange', 'topic': 'eqiad.mediawiki.recentchange', 'uri': 'https://commons.wikimedia.org/wiki/Category:Iron_Age_in_Norway'},
'namespace': 14,
'parsedcomment': '<a href="/wiki/...">...</a> removed from category',
'server_name': 'commons.wikimedia.org',
'server_script_path': '/w',
'server_url': 'https://commons.wikimedia.org',
'timestamp': 1652349420,
'title': 'Category:Iron Age in Norway',
'type': 'categorize',
'user': 'Krg',
'wiki': 'commonswiki'}
(각 이벤트는 최근 변경 스트림의 개별 이벤트로, meta, title, user, wiki, comment 등의 필드를 담고 있어요.)
최근 변경 사항을 Kafka로 수집
이제 각 이벤트를 Apache Kafka로 가져올 거예요. 먼저 5개 파티션을 가진 wiki_events라는 Kafka 토픽을 만들어요:
docker exec -it kafka-wiki /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 \
--create \
--topic wiki_events \
--partitions 5
wiki_to_kafka.py라는 새 파일을 만들고 다음 라이브러리를 가져와요:
import json
import sseclient
import datetime
import requests
import time
from confluent_kafka import Producer
wiki_to_kafka.py
다음 함수들을 추가하세요:
def with_requests(url, headers):
"""Get a streaming response for the given event feed using requests."""
return requests.get(url, stream=True, headers=headers)
def acked(err, msg):
if err is not None:
print("Failed to deliver message: {0}: {1}"
.format(msg.value(), err.str()))
def json_serializer(obj):
if isinstance(obj, (datetime.datetime, datetime.date)):
return obj.isoformat()
raise "Type %s not serializable" % type(obj)
wiki_to_kafka.py
이제 최근 변경 API를 호출하고 wiki_events 토픽으로 이벤트를 가져오는 코드를 추가해요:
producer = Producer({'bootstrap.servers': 'localhost:9092'})
url = 'https://stream.wikimedia.org/v2/stream/recentchange'
headers = {'Accept': 'text/event-stream'}
response = with_requests(url, headers)
client = sseclient.SSEClient(response)
events_processed = 0
while True:
try:
for event in client.events():
stream = json.loads(event.data)
payload = json.dumps(stream, default=json_serializer, ensure_ascii=False).encode('utf-8')
producer.produce(topic='wiki_events',
key=str(stream['meta']['id']), value=payload, callback=acked)
events_processed += 1
if events_processed % 100 == 0:
print(f"{str(datetime.datetime.now())} Flushing after {events_processed} events")
producer.flush()
except Exception as ex:
print(f"{str(datetime.datetime.now())} Got error:" + str(ex))
response = with_requests(url, headers)
client = sseclient.SSEClient(response)
time.sleep(2)
wiki_to_kafka.py
이 스크립트의 강조된 부분은 이벤트가 Kafka로 수집되고 디스크로 flush 되는 위치를 나타내요.
이 스크립트를 실행하면:
python wiki_to_kafka.py
100개 메시지가 Kafka에 push될 때마다 메시지를 볼 거예요:
출력
2022-05-12 10:58:34.449326 Flushing after 100 events
2022-05-12 10:58:39.151599 Flushing after 200 events
2022-05-12 10:58:43.399528 Flushing after 300 events
...
Kafka 살펴보기
데이터가 Kafka에 들어갔는지 확인해보아요.
다음 명령은 wiki_events 토픽의 각 파티션에 대한 메시지 오프셋을 반환해요:
docker exec -it kafka-wiki /opt/kafka/bin/kafka-get-offsets.sh \
--bootstrap-server localhost:9092 \
--topic wiki_events
출력
wiki_events:0:42
wiki_events:1:61
wiki_events:2:52
wiki_events:3:56
wiki_events:4:58
좋아 보여요. 다음 명령을 실행해 이 토픽의 모든 메시지를 스트리밍할 수도 있어요:
docker exec -it kafka-wiki /opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server localhost:9092 \
--topic wiki_events \
--from-beginning
출력
...
{"$schema": "/mediawiki/recentchange/1.0.0", "meta": {"uri": "https://en.wikipedia.org/wiki/Super_Wings", ...}, "type": "log", ...}
{"$schema": "/mediawiki/recentchange/1.0.0", "meta": {"uri": "https://no.wikipedia.org/wiki/Brukerdiskusjon:Haros", ...}, "id": 84572581, "type": "edit", ...}
^CProcessed a total of 269 messages
Pinot 구성
이제 Kafka의 데이터를 소비하도록 Pinot을 구성해보아요.
다음 스키마를 만들 거예요:
{
"schemaName": "wikipedia",
"dimensionFieldSpecs": [
{ "name": "id", "dataType": "STRING" },
{ "name": "wiki", "dataType": "STRING" },
{ "name": "user", "dataType": "STRING" },
{ "name": "title", "dataType": "STRING" },
{ "name": "comment", "dataType": "STRING" },
{ "name": "stream", "dataType": "STRING" },
{ "name": "domain", "dataType": "STRING" },
{ "name": "topic", "dataType": "STRING" },
{ "name": "type", "dataType": "STRING" },
{ "name": "uri", "dataType": "STRING" },
{ "name": "bot", "dataType": "BOOLEAN" },
{ "name": "metaJson", "dataType": "STRING" }
],
"dateTimeFieldSpecs": [
{
"name": "ts",
"dataType": "TIMESTAMP",
"format": "1:MILLISECONDS:EPOCH",
"granularity": "1:MILLISECONDS"
}
]
}
schema.json
그리고 다음 테이블 구성을 사용해요:
{
"tableName": "wikievents",
"tableType": "REALTIME",
"segmentsConfig": {
"timeColumnName": "ts",
"schemaName": "wikipedia",
"replication": "1",
"replicasPerPartition": "1"
},
"tableIndexConfig": {
"invertedIndexColumns": [],
"rangeIndexColumns": [],
"autoGeneratedInvertedIndex": false,
"createInvertedIndexDuringSegmentGeneration": false,
"sortedColumn": [],
"bloomFilterColumns": [],
"loadMode": "MMAP",
"streamConfigs": {
"streamType": "kafka",
"stream.kafka.topic.name": "wiki_events",
"stream.kafka.broker.list": "kafka-wiki:9093",
"stream.kafka.consumer.prop.auto.offset.reset": "smallest",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.json.JSONMessageDecoder",
"realtime.segment.flush.threshold.rows": "1000",
"realtime.segment.flush.threshold.time": "24h",
"realtime.segment.flush.segment.size": "100M"
},
"tenants": {
"broker": "DefaultTenant",
"server": "DefaultTenant",
"tagOverrideConfig": {}
},
"noDictionaryColumns": [],
"onHeapDictionaryColumns": [],
"varLengthDictionaryColumns": [],
"enableDefaultStarTree": false,
"enableDynamicStarTreeCreation": false,
"aggregateMetrics": false,
"nullHandlingEnabled": false
},
"metadata": {},
"quota": {},
"routing": {},
"query": {},
"ingestionConfig": {
"transformConfigs": [
{ "columnName": "metaJson", "transformFunction": "JSONFORMAT(meta)" },
{ "columnName": "id", "transformFunction": "JSONPATH(metaJson, '$.id')" },
{ "columnName": "stream", "transformFunction": "JSONPATH(metaJson, '$.stream')" },
{ "columnName": "domain", "transformFunction": "JSONPATH(metaJson, '$.domain')" },
{ "columnName": "topic", "transformFunction": "JSONPATH(metaJson, '$.topic')" },
{ "columnName": "uri", "transformFunction": "JSONPATH(metaJson, '$.uri')" },
{ "columnName": "ts", "transformFunction": "\"timestamp\" * 1000" }
]
},
"isDimTable": false
}
table.json
강조된 줄은 이벤트가 들어 있는 Kafka 토픽에 Pinot을 연결하는 방법이에요. 다음 명령으로 스키마와 테이블을 만들세요:
docker exec -it pinot-controller-wiki bin/pinot-admin.sh AddTable \
-tableConfigFile /config/table.json \
-schemaFile /config/schema.json \
-exec
완료한 후 Pinot UI로 이동해 다음 쿼리를 실행해 데이터가 Pinot에 들어갔는지 확인하세요:
select domain, count(*)
from wikievents
group by domain
order by count(*) DESC
limit 10
레코드가 보이면 모든 것이 예상대로 동작하고 있는 거예요.
Dash 대시보드 만들기
이제 Pinot에 몇 가지 더 많은 쿼리를 작성하고 결과를 Dash에 표시해보아요.
먼저 다음 라이브러리를 설치하세요:
pip install dash pinotdb plotly pandas
dashboard.py라는 파일을 만들어 라이브러리를 가져오고 페이지의 헤더를 작성해요:
import pandas as pd
from dash import Dash, html, dcc
import plotly.graph_objects as go
from pinotdb import connect
import plotly.express as px
external_stylesheets = ['https://codepen.io/chriddyp/pen/bWLwgP.css']
app = Dash(__name__, external_stylesheets=external_stylesheets)
app.title = "Wiki Recent Changes Dashboard"
app.py
Pinot에 연결하고 최근 변경 사항, 변경한 사용자, 변경이 발생한 도메인을 반환하는 쿼리를 작성해요:
conn = connect(host='localhost', port=8099, path='/query/sql', scheme='http')
query = """select
count(*) FILTER(WHERE ts > ago('PT1M')) AS events1Min,
count(*) FILTER(WHERE ts <= ago('PT1M') AND ts > ago('PT2M')) AS events1Min2Min,
distinctcount(user) FILTER(WHERE ts > ago('PT1M')) AS users1Min,
distinctcount(user) FILTER(WHERE ts <= ago('PT1M') AND ts > ago('PT2M')) AS users1Min2Min,
distinctcount(domain) FILTER(WHERE ts > ago('PT1M')) AS domains1Min,
distinctcount(domain) FILTER(WHERE ts <= ago('PT1M') AND ts > ago('PT2M')) AS domains1Min2Min
from wikievents
where ts > ago('PT2M')
limit 1
"""
curs = conn.cursor()
curs.execute(query)
df_summary = pd.DataFrame(curs, columns=[item[0] for item in curs.description])
app.py
쿼리의 강조된 부분은 지난 1분과 그 이전 1분의 이벤트 수를 세는 방법을 보여줘요. 그런 다음 유사한 방식으로 고유 사용자와 도메인 수를 셉니다.
메트릭 (Metrics)
이제 해당 데이터를 기반으로 몇 가지 메트릭을 만들어요.
먼저 이러한 메트릭을 만들기 위한 몇 가지 헬퍼 함수를 만들어요:
from dash import html, dash_table
import plotly.graph_objects as go
def add_delta_trace(fig, title, value, last_value, row, column):
fig.add_trace(go.Indicator(
mode = "number+delta",
title= {'text': title},
value = value,
delta = {'reference': last_value, 'relative': True},
domain = {'row': row, 'column': column})
)
def add_trace(fig, title, value, row, column):
fig.add_trace(go.Indicator(
mode = "number",
title= {'text': title},
value = value,
domain = {'row': row, 'column': column})
)
dash_utils.py
이제 app.py에 다음 import를 추가해요:
from dash_utils import add_delta_trace, add_trace
app.py
그리고 파일 끝에 다음 코드를 추가해요:
fig = go.Figure(layout=go.Layout(height=300))
if df_summary["events1Min"][0] > 0:
if df_summary["events1Min"][0] > 0:
add_delta_trace(fig, "Changes", df_summary["events1Min"][0], df_summary["events1Min2Min"][0], 0, 0)
add_delta_trace(fig, "Users", df_summary["users1Min"][0], df_summary["users1Min2Min"][0], 0, 1)
add_delta_trace(fig, "Domain", df_summary["domains1Min"][0], df_summary["domains1Min2Min"][0], 0, 2)
else:
add_trace(fig, "Changes", df_summary["events1Min"][0], 0, 0)
add_trace(fig, "Users", df_summary["users1Min2Min"][0], 0, 1)
add_trace(fig, "Domains", df_summary["domains1Min2Min"][0], 0, 2)
fig.update_layout(grid = {"rows": 1, "columns": 3, 'pattern': "independent"},)
else:
fig.update_layout(annotations = [{"text": "No events found", "xref": "paper", "yref": "paper", "showarrow": False, "font": {"size": 28}}])
app.layout = html.Div([
html.H1("Wiki Recent Changes Dashboard", style={'text-align': 'center'}),
html.Div(id='content', children=[
dcc.Graph(figure=fig)
])
])
if __name__ == '__main__':
app.run_server(debug=True)
app.py
터미널로 돌아가 다음 명령을 실행해요:
python dashboard.py
localhost:8051로 이동하면 Dash 앱을 볼 수 있어요.
분당 변경 수 (Changes per minute)
다음으로 Wikimedia에 분당 가해지는 변경 수를 보여주는 선 차트를 추가해요. app.py를 다음과 같이 업데이트하세요:
query = """
select ToDateTime(DATETRUNC('minute', ts), 'yyyy-MM-dd hh:mm:ss') AS dateMin, count(*) AS changes,
distinctcount(user) AS users,
distinctcount(domain) AS domains
from wikievents
where ts > ago('PT2M')
group by dateMin
order by dateMin desc
LIMIT 30
"""
curs.execute(query)
df_ts = pd.DataFrame(curs, columns=[item[0] for item in curs.description])
df_ts_melt = pd.melt(df_ts, id_vars=['dateMin'], value_vars=['changes', 'users', 'domains'])
line_chart = px.line(df_ts_melt, x='dateMin', y="value", color='variable', color_discrete_sequence =['blue', 'red', 'green'])
line_chart['layout'].update(margin=dict(l=0,r=0,b=0,t=40), title="Changes/Users/Domains per minute")
line_chart.update_yaxes(range=[0, df_ts["changes"].max() * 1.1])
app.layout = html.Div([
html.H1("Wiki Recent Changes Dashboard", style={'text-align': 'center'}),
html.Div(id='content', children=[
dcc.Graph(figure=fig),
dcc.Graph(figure=line_chart),
])
])
app.py
웹 브라우저로 돌아가면 다음과 같은 것을 볼 수 있어요.
자동 새로고침 (Auto Refresh)
지금은 웹 브라우저를 새로고침해 메트릭과 선 차트를 업데이트해야 하지만, 그 과정이 자동으로 이뤄지면 훨씬 낫겠죠. 자동 새로고침 기능을 추가해보아요.
이를 위해서는 각 구성요소가 콜백으로 주석 처리된 함수에서 렌더링되고, 그 함수가 간격에 따라 호출되도록 애플리케이션을 재구성해야 해요.
앱 레이아웃은 이제 다음과 같아요:
app.layout = html.Div([
html.H1("Wiki Recent Changes Dashboard", style={'text-align': 'center'}),
html.Div(id='latest-timestamp', style={"padding": "5px 0", "text-align": "center"}),
dcc.Interval(
id='interval-component',
interval=1 * 1000,
n_intervals=0
),
html.Div(id='content', children=[
dcc.Graph(id="indicators"),
dcc.Graph(id="time-series"),
])
])
app.py
interval-component는 1,000밀리초마다 콜백을 실행하도록 구성돼요.latest-timestamp는 최신 타임스탬프를 담는 컨테이너예요.indicators는 사용자, 도메인, 변경 수의 최신 개수를 담는 인디케이터를 포함해요.time-series는 시계열 선 차트를 포함해요.
타임스탬프는 다음 콜백 함수로 새로고침돼요:
@app.callback(
Output(component_id='latest-timestamp', component_property='children'),
Input('interval-component', 'n_intervals'))
def timestamp(n):
return html.Span(f"Last updated: {datetime.datetime.now()}")
app.py
인디케이터는 다음 함수로 새로고침돼요:
@app.callback(Output(component_id='indicators', component_property='figure'),
Input('interval-component', 'n_intervals'))
def indicators(n):
query = """
select count(*) FILTER(WHERE ts > ago('PT1M')) AS events1Min,
count(*) FILTER(WHERE ts <= ago('PT1M') AND ts > ago('PT2M')) AS events1Min2Min,
distinctcount(user) FILTER(WHERE ts > ago('PT1M')) AS users1Min,
distinctcount(user) FILTER(WHERE ts <= ago('PT1M') AND ts > ago('PT2M')) AS users1Min2Min,
distinctcount(domain) FILTER(WHERE ts > ago('PT1M')) AS domains1Min,
distinctcount(domain) FILTER(WHERE ts <= ago('PT1M') AND ts > ago('PT2M')) AS domains1Min2Min
from wikievents
where ts > ago('PT2M')
limit 1
"""
curs = connection.cursor()
curs.execute(query)
df_summary = pd.DataFrame(curs, columns=[item[0] for item in curs.description])
curs.close()
fig = go.Figure(layout=go.Layout(height=300))
if df_summary["events1Min"][0] > 0:
if df_summary["events1Min"][0] > 0:
add_delta_trace(fig, "Changes", df_summary["events1Min"][0], df_summary["events1Min2Min"][0], 0, 0)
add_delta_trace(fig, "Users", df_summary["users1Min"][0], df_summary["users1Min2Min"][0], 0, 1)
add_delta_trace(fig, "Domain", df_summary["domains1Min"][0], df_summary["domains1Min2Min"][0], 0, 2)
else:
add_trace(fig, "Changes", df_summary["events1Min"][0], 0, 0)
add_trace(fig, "Users", df_summary["users1Min2Min"][0], 0, 1)
add_trace(fig, "Domains", df_summary["domains1Min2Min"][0], 0, 2)
fig.update_layout(grid = {"rows": 1, "columns": 3, 'pattern': "independent"},)
else:
fig.update_layout(annotations = [{"text": "No events found", "xref": "paper", "yref": "paper", "showarrow": False, "font": {"size": 28}}])
return fig
app.py
마지막으로 다음 함수가 선 차트를 새로고침해요:
@app.callback(Output(component_id='time-series', component_property='figure'),
Input('interval-component', 'n_intervals'))
def time_series(n):
query = """
select ToDateTime(DATETRUNC('minute', ts), 'yyyy-MM-dd hh:mm:ss') AS dateMin, count(*) AS changes,
distinctcount(user) AS users,
distinctcount(domain) AS domains
from wikievents
where ts > ago('PT1H')
group by dateMin
order by dateMin desc
LIMIT 30
"""
curs = connection.cursor()
curs.execute(query)
df_ts = pd.DataFrame(curs, columns=[item[0] for item in curs.description])
curs.close()
df_ts_melt = pd.melt(df_ts, id_vars=['dateMin'], value_vars=['changes', 'users', 'domains'])
line_chart = px.line(df_ts_melt, x='dateMin', y="value", color='variable', color_discrete_sequence =['blue', 'red', 'green'])
line_chart['layout'].update(margin=dict(l=0,r=0,b=0,t=40), title="Changes/Users/Domains per minute")
line_chart.update_yaxes(range=[0, df_ts["changes"].max() * 1.1])
return line_chart
app.py
웹 브라우저로 돌아가면 다음을 볼 수 있어요.
이 예시에서 사용된 전체 스크립트는 아래와 같아요:
import pandas as pd
from dash import Dash, html, dash_table, dcc, Input, Output
import plotly.graph_objects as go
from pinotdb import connect
from dash_utils import add_delta_trace, add_trace
import plotly.express as px
import datetime
external_stylesheets = ['https://codepen.io/chriddyp/pen/bWLwgP.css']
app = Dash(__name__, external_stylesheets=external_stylesheets)
app.title = "Wiki Recent Changes Dashboard"
connection = connect(host="localhost", port="8099", path="/query/sql", scheme=( "http"))
@app.callback(Output(component_id='indicators', component_property='figure'),
Input('interval-component', 'n_intervals'))
def indicators(n):
query = """
select count(*) FILTER(WHERE ts > ago('PT1M')) AS events1Min,
count(*) FILTER(WHERE ts <= ago('PT1M') AND ts > ago('PT2M')) AS events1Min2Min,
distinctcount(user) FILTER(WHERE ts > ago('PT1M')) AS users1Min,
distinctcount(user) FILTER(WHERE ts <= ago('PT1M') AND ts > ago('PT2M')) AS users1Min2Min,
distinctcount(domain) FILTER(WHERE ts > ago('PT1M')) AS domains1Min,
distinctcount(domain) FILTER(WHERE ts <= ago('PT1M') AND ts > ago('PT2M')) AS domains1Min2Min
from wikievents
where ts > ago('PT2M')
limit 1
"""
curs = connection.cursor()
curs.execute(query)
df_summary = pd.DataFrame(curs, columns=[item[0] for item in curs.description])
curs.close()
fig = go.Figure(layout=go.Layout(height=300))
if df_summary["events1Min"][0] > 0:
if df_summary["events1Min"][0] > 0:
add_delta_trace(fig, "Changes", df_summary["events1Min"][0], df_summary["events1Min2Min"][0], 0, 0)
add_delta_trace(fig, "Users", df_summary["users1Min"][0], df_summary["users1Min2Min"][0], 0, 1)
add_delta_trace(fig, "Domain", df_summary["domains1Min"][0], df_summary["domains1Min2Min"][0], 0, 2)
else:
add_trace(fig, "Changes", df_summary["events1Min"][0], 0, 0)
add_trace(fig, "Users", df_summary["users1Min2Min"][0], 0, 1)
add_trace(fig, "Domains", df_summary["domains1Min2Min"][0], 0, 2)
fig.update_layout(grid = {"rows": 1, "columns": 3, 'pattern': "independent"},)
else:
fig.update_layout(annotations = [{"text": "No events found", "xref": "paper", "yref": "paper", "showarrow": False, "font": {"size": 28}}])
return fig
@app.callback(Output(component_id='time-series', component_property='figure'),
Input('interval-component', 'n_intervals'))
def time_series(n):
query = """
select ToDateTime(DATETRUNC('minute', ts), 'yyyy-MM-dd hh:mm:ss') AS dateMin, count(*) AS changes,
distinctcount(user) AS users,
distinctcount(domain) AS domains
from wikievents
where ts > ago('PT1H')
group by dateMin
order by dateMin desc
LIMIT 30
"""
curs = connection.cursor()
curs.execute(query)
df_ts = pd.DataFrame(curs, columns=[item[0] for item in curs.description])
curs.close()
df_ts_melt = pd.melt(df_ts, id_vars=['dateMin'], value_vars=['changes', 'users', 'domains'])
line_chart = px.line(df_ts_melt, x='dateMin', y="value", color='variable', color_discrete_sequence =['blue', 'red', 'green'])
line_chart['layout'].update(margin=dict(l=0,r=0,b=0,t=40), title="Changes/Users/Domains per minute")
line_chart.update_yaxes(range=[0, df_ts["changes"].max() * 1.1])
return line_chart
@app.callback(
Output(component_id='latest-timestamp', component_property='children'),
Input('interval-component', 'n_intervals'))
def timestamp(n):
return html.Span(f"Last updated: {datetime.datetime.now()}")
app.layout = html.Div([
html.H1("Wiki Recent Changes Dashboard", style={'text-align': 'center'}),
html.Div(id='latest-timestamp', style={"padding": "5px 0", "text-align": "center"}),
dcc.Interval(
id='interval-component',
interval=1 * 1000,
n_intervals=0
),
html.Div(id='content', children=[
dcc.Graph(id="indicators"),
dcc.Graph(id="time-series"),
])
])
if __name__ == '__main__':
app.run_server(debug=True)
dashboard.py
요약 (Summary)
이 가이드에서는 Wikimedia 이벤트 스트림에서 데이터를 Kafka로 게시하고, 거기서 Pinot으로 수집하고, 마지막으로 Dash에서 실행되는 SQL 쿼리를 사용해 데이터를 이해하는 방법을 배웠어요.