|
48 | 48 | "id": "08447511-b105-416f-a959-285330d417ba", |
49 | 49 | "metadata": {}, |
50 | 50 | "source": [ |
51 | | - "**Note**: This tutorial has been tested for Python 3.11, using JupyterLab 4.2.0 and perspective-python 2.10.0." |
| 51 | + "**Note**: This tutorial has been tested for Python 3.11, using JupyterLab 4.2.0 and perspective-python 3.1.5." |
52 | 52 | ] |
53 | 53 | }, |
54 | 54 | { |
|
130 | 130 | "metadata": {}, |
131 | 131 | "outputs": [], |
132 | 132 | "source": [ |
133 | | - "from perspective import PerspectiveWidget, Plugin\n", |
| 133 | + "from perspective.widget import PerspectiveWidget\n", |
134 | 134 | "\n", |
135 | 135 | "# Data schema\n", |
136 | | - "data = {\"servername\": str, \"timestamp\": datetime, \"event\": str, \"servername_count\": int}\n", |
137 | | - "widget = PerspectiveWidget(data, plugin=Plugin.XBAR, group_by=[\"servername\"], columns=[\"servername_count\"], theme='Pro Light')\n", |
| 136 | + "data = {\"servername\": \"string\", \"timestamp\": \"datetime\", \"event\": \"string\", \"servername_count\": \"integer\"}\n", |
| 137 | + "#widget = PerspectiveWidget(data, plugin=\"X Bar\", group_by=[\"servername\"], columns=[\"servername_count\"], aggregates={\"servername_count\": \"last\"}, theme='Pro Light', binding_mode=\"client-server\")\n", |
| 138 | + "widget = PerspectiveWidget(data, plugin=\"X Bar\", group_by=[\"servername\"], columns=[\"servername_count\"], aggregates={\"servername_count\": \"last\"}, theme='Pro Light', binding_mode=\"client-server\")\n", |
138 | 139 | "widget" |
139 | 140 | ] |
140 | 141 | }, |
|
213 | 214 | " if self._running:\n", |
214 | 215 | " self._running = False\n", |
215 | 216 | " self._thread.join()\n", |
216 | | - " self._source.close()\n", |
| 217 | + " self._source.resp.close()\n", |
217 | 218 | "\n", |
218 | 219 | " def _run(self):\n", |
219 | 220 | " servernames = dict([])\n", |
220 | | - " while self._running:\n", |
221 | | - " for item in self._source:\n", |
222 | | - " if item.event == 'message':\n", |
223 | | - " try:\n", |
224 | | - " change = json.loads(item.data)\n", |
225 | | - " except ValueError:\n", |
226 | | - " pass\n", |
227 | | - " else:\n", |
228 | | - " # discard canary events\n", |
229 | | - " # WMF Data Engineering team produces artificial 'canary' events into \n", |
230 | | - " # each stream multiple times an hour. The presence of these canary\n", |
231 | | - " # events in a stream allow us to differentiate between a broken event\n", |
232 | | - " # stream, and an empty one.\n", |
233 | | - " # We will also filter bot-generated events\n", |
234 | | - " if change['meta']['domain'] == 'canary' or re.search('bot', change['user'], re.IGNORECASE):\n", |
235 | | - " continue\n", |
236 | | - " timestamp = change['meta']['dt']\n", |
237 | | - " event = f\"{timestamp}:: {change['user']} edited {change['title']}\"\n", |
238 | | - " servername = change['server_name']\n", |
239 | | - " # Manually \"tick\" this edge\n", |
240 | | - " self.push_tick(\n", |
241 | | - " WikiData(\n", |
242 | | - " servername=servername,\n", |
243 | | - " servername_count=0,\n", |
244 | | - " timestamp=timestamp,\n", |
245 | | - " event=event\n", |
246 | | - " )\n", |
| 221 | + " for item in self._source:\n", |
| 222 | + " if not self._running:\n", |
| 223 | + " break\n", |
| 224 | + " if item.event == 'message':\n", |
| 225 | + " try:\n", |
| 226 | + " change = json.loads(item.data)\n", |
| 227 | + " except ValueError:\n", |
| 228 | + " pass\n", |
| 229 | + " else:\n", |
| 230 | + " # discard canary events\n", |
| 231 | + " # WMF Data Engineering team produces artificial 'canary' events into \n", |
| 232 | + " # each stream multiple times an hour. The presence of these canary\n", |
| 233 | + " # events in a stream allow us to differentiate between a broken event\n", |
| 234 | + " # stream, and an empty one.\n", |
| 235 | + " # We will also filter bot-generated events\n", |
| 236 | + " if change['meta']['domain'] == 'canary' or re.search('bot', change['user'], re.IGNORECASE):\n", |
| 237 | + " continue\n", |
| 238 | + " timestamp = change['meta']['dt']\n", |
| 239 | + " event = f\"{timestamp}:: {change['user']} edited {change['title']}\"\n", |
| 240 | + " servername = change['server_name']\n", |
| 241 | + " # Manually \"tick\" this edge\n", |
| 242 | + " self.push_tick(\n", |
| 243 | + " WikiData(\n", |
| 244 | + " servername=servername,\n", |
| 245 | + " servername_count=0,\n", |
| 246 | + " timestamp=timestamp,\n", |
| 247 | + " event=event\n", |
247 | 248 | " )\n", |
| 249 | + " )\n", |
248 | 250 | "\n", |
249 | 251 | "# Create the graph-time representation of our adapter\n", |
250 | 252 | "FetchWikiData = py_push_adapter_def(\"FetchWikiData\", FetchWikiDataAdapter, csp.ts[WikiData], url=str)" |
|
268 | 270 | "outputs": [], |
269 | 271 | "source": [ |
270 | 272 | "@csp.node\n", |
271 | | - "def update_widget(wiki_event: csp.ts[WikiData], widget: PerspectiveWidget):\n", |
| 273 | + "def update_widget(wiki_event: csp.ts[WikiData], widget: PerspectiveWidget, throttle: timedelta = timedelta(seconds=0.5)):\n", |
| 274 | + " # Updates the perspective widget with batched updates for scalability\n", |
| 275 | + " with csp.alarms():\n", |
| 276 | + " alarm = csp.alarm(bool)\n", |
| 277 | + "\n", |
| 278 | + " with csp.state():\n", |
| 279 | + " s_buffer = []\n", |
| 280 | + "\n", |
| 281 | + " with csp.start():\n", |
| 282 | + " csp.schedule_alarm(alarm, throttle, True)\n", |
| 283 | + " \n", |
272 | 284 | " if csp.ticked(wiki_event):\n", |
273 | | - " widget.update({\n", |
274 | | - " \"servername\": [wiki_event.servername],\n", |
275 | | - " \"servername_count\": [wiki_event.servername_count],\n", |
276 | | - " \"timestamp\": [wiki_event.timestamp],\n", |
277 | | - " \"event\": [wiki_event.event],\n", |
| 285 | + " s_buffer.append({\n", |
| 286 | + " \"servername\": wiki_event.servername,\n", |
| 287 | + " \"servername_count\": wiki_event.servername_count,\n", |
| 288 | + " \"timestamp\": wiki_event.timestamp,\n", |
| 289 | + " \"event\": wiki_event.event,\n", |
278 | 290 | " })\n", |
279 | 291 | "\n", |
| 292 | + " if csp.ticked(alarm):\n", |
| 293 | + " if len(s_buffer) > 0:\n", |
| 294 | + " widget.update(s_buffer)\n", |
| 295 | + " s_buffer = []\n", |
| 296 | + "\n", |
| 297 | + " csp.schedule_alarm(alarm, throttle, True)\n", |
| 298 | + "\n", |
280 | 299 | "@csp.node\n", |
281 | 300 | "def compute_server_count(wiki_event: csp.ts[WikiData]) -> csp.ts[WikiData]:\n", |
282 | 301 | " # takes the raw struct in, creates a copy with the count set and ticks it out\n", |
|
477 | 496 | "@csp.graph\n", |
478 | 497 | "def wiki_graph():\n", |
479 | 498 | " print(\"Start of graph building\")\n", |
480 | | - " if csp.is_configured_realtime:\n", |
| 499 | + " if csp.is_configured_realtime():\n", |
481 | 500 | " URL = \"https://stream.wikimedia.org/v2/stream/recentchange\"\n", |
482 | 501 | " events = FetchWikiData(url=URL)\n", |
483 | 502 | " else:\n", |
|
540 | 559 | "name": "python", |
541 | 560 | "nbconvert_exporter": "python", |
542 | 561 | "pygments_lexer": "ipython3", |
543 | | - "version": "3.11.9" |
| 562 | + "version": "3.11.7" |
544 | 563 | } |
545 | 564 | }, |
546 | 565 | "nbformat": 4, |
|
0 commit comments