Streamlit에 연결하기
Streamlit에 연결하기 (Connect to Streamlit)
이 Apache Pinot 가이드에서는 Streamlit 웹 프레임워크를 사용해 데이터를 시각화하는 방법을 알려드려요.
이 가이드에서는 Streamlit을 사용해 Apache Pinot의 데이터를 시각화하는 방법을 배워요. Streamlit은 대화형 데이터 기반 웹 앱을 쉽게 만들게 해주는 Python 라이브러리예요.
Streamlit을 사용해 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 최근 변경 스트림
Streamlit 가이드의 Dash 예시와 동일한 Wikimedia SSE 스트림을 사용해요. 엔드포인트는 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 클라이언트 라이브러리를 사용해 최근 변경 피드에 연결하는 방법을 보여줘요.
이 스크립트를 실행하면 (잘린) 다음 출력을 볼 거예요:
출력
{'$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'}
최근 변경 사항을 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 되는 위치를 나타내요.
이 스크립트를 실행하면 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
...
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
레코드가 보이면 모든 것이 예상대로 동작하고 있는 거예요.
Streamlit 대시보드 만들기
이제 Pinot에 몇 가지 더 많은 쿼리를 작성하고 결과를 Streamlit에 표시해보아요.
먼저 다음 라이브러리를 설치하세요:
pip install streamlit pinotdb plotly pandas
app.py라는 파일을 만들어 라이브러리를 가져오고 페이지의 헤더를 작성해요:
import pandas as pd
import streamlit as st
from pinotdb import connect
import plotly.express as px
st.set_page_config(layout="wide")
st.header("Wikipedia Recent Changes")
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)
이제 해당 데이터를 기반으로 몇 가지 메트릭을 만들어요:
metric1, metric2, metric3 = st.columns(3)
metric1.metric(label="Changes", value=df_summary['events1Min'].values[0],
delta=float(df_summary['events1Min'].values[0] - df_summary['events1Min2Min'].values[0]))
metric2.metric(label="Users", value=df_summary['users1Min'].values[0],
delta=float(df_summary['users1Min'].values[0] - df_summary['users1Min2Min'].values[0]))
metric3.metric(label="Domains", value=df_summary['domains1Min'].values[0],
delta=float(df_summary['domains1Min'].values[0] - df_summary['domains1Min2Min'].values[0]))
app.py
터미널로 돌아가 다음 명령을 실행해요:
streamlit run app.py
localhost:8501로 이동하면 Streamlit 앱을 볼 수 있어요.
분당 변경 수 (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('PT1H')
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'])
fig = px.line(df_ts_melt, x='dateMin', y="value", color='variable', color_discrete_sequence =['blue', 'red', 'green'])
fig['layout'].update(margin=dict(l=0,r=0,b=0,t=40), title="Changes/Users/Domains per minute")
fig.update_yaxes(range=[0, df_ts["changes"].max() * 1.1])
st.plotly_chart(fig, use_container_width=True)
app.py
웹 브라우저로 돌아가면 다음과 같은 것을 볼 수 있어요.
자동 새로고침 (Auto Refresh)
지금은 웹 브라우저를 새로고침해 메트릭과 선 차트를 업데이트해야 하지만, 그 과정이 자동으로 이뤄지면 훨씬 낫겠죠. 자동 새로고침 기능을 추가해보아요.
app.py 맨 위의 헤더 바로 아래에 다음 코드를 추가하세요:
if not "sleep_time" in st.session_state:
st.session_state.sleep_time = 2
if not "auto_refresh" in st.session_state:
st.session_state.auto_refresh = True
auto_refresh = st.checkbox('Auto Refresh?', st.session_state.auto_refresh)
if auto_refresh:
number = st.number_input('Refresh rate in seconds', value=st.session_state.sleep_time)
st.session_state.sleep_time = number
app.py
그리고 파일 맨 끝에 다음 코드를 추가해요:
if auto_refresh:
time.sleep(number)
st.experimental_rerun()
app.py
웹 브라우저로 돌아가면 다음을 볼 수 있어요.
이 예시에서 사용된 전체 스크립트는 아래와 같아요:
import pandas as pd
import streamlit as st
from pinotdb import connect
from datetime import datetime
import plotly.express as px
import time
st.set_page_config(layout="wide")
conn = connect(host='localhost', port=8099, path='/query/sql', scheme='http')
st.header("Wikipedia Recent Changes")
now = datetime.now()
dt_string = now.strftime("%d %B %Y %H:%M:%S")
st.write(f"Last update: {dt_string}")
# Use session state to keep track of whether we need to auto refresh the page and the refresh frequency
if not "sleep_time" in st.session_state:
st.session_state.sleep_time = 2
if not "auto_refresh" in st.session_state:
st.session_state.auto_refresh = True
auto_refresh = st.checkbox('Auto Refresh?', st.session_state.auto_refresh)
if auto_refresh:
number = st.number_input('Refresh rate in seconds', value=st.session_state.sleep_time)
st.session_state.sleep_time = number
# 지난 1분 동안 발생한 변경 찾기
# 1분에서 2분 사이에 발생한 변경 찾기
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])
metric1, metric2, metric3 = st.columns(3)
metric1.metric(
label="Changes",
value=df_summary['events1Min'].values[0],
delta=float(df_summary['events1Min'].values[0] - df_summary['events1Min2Min'].values[0])
)
metric2.metric(
label="Users",
value=df_summary['users1Min'].values[0],
delta=float(df_summary['users1Min'].values[0] - df_summary['users1Min2Min'].values[0])
)
metric3.metric(
label="Domains",
value=df_summary['domains1Min'].values[0],
delta=float(df_summary['domains1Min'].values[0] - df_summary['domains1Min2Min'].values[0])
)
# 지난 1시간 동안 분(minute)별 모든 변경 찾기
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('PT10M')
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'])
fig = px.line(df_ts_melt, x='dateMin', y="value", color='variable', color_discrete_sequence =['blue', 'red', 'green'])
fig['layout'].update(margin=dict(l=0,r=0,b=0,t=40), title="Changes/Users/Domains per minute")
fig.update_yaxes(range=[0, df_ts["changes"].max() * 1.1])
st.plotly_chart(fig, use_container_width=True)
# 페이지 새로고침
if auto_refresh:
time.sleep(number)
st.experimental_rerun()
app.py
요약 (Summary)
이 가이드에서는 Wikimedia 이벤트 스트림에서 데이터를 Kafka로 게시하고, 거기서 Pinot으로 수집하고, 마지막으로 Streamlit에서 실행되는 SQL 쿼리를 사용해 데이터를 이해하는 방법을 배웠어요.