Skip to content

BGPStream Class

pybgpflux.bgpstream.BGPStream

Stream and process BGP messages from multiple collectors.

BGPStream is a high-performance alternative to PyBGPStream that parses BGP MRT files using BGPKIT. It can stream both historical and live BGP data with support for advanced filtering, multiple parser backends, and memory-efficient lazy loading.

Attributes:

Name Type Description
collectors list[str]

List of collector names to fetch data from.

data_type list[Literal['ribs', 'updates']]

Data types to stream ("ribs" or "updates").

ts_start float | None

Start timestamp (Unix epoch). None for live mode.

ts_end float | None

End timestamp (Unix epoch). None for live mode.

filters FilterOptions

Filtering options for BGP elements.

cache_dir Directory | TemporaryDirectory

Cache directory for downloaded files.

parser_name str

Backend parser to use ("pybgpkit", "bgpkit", "bgpdump", "pybgpstream").

max_concurrent_downloads int

Maximum concurrent file downloads.

ram_fetch bool

Use RAM disk (/dev/shm, /Volumes/RAMDisk) if available.

jitter_buffer_delay float

Delay (seconds) for jitter buffer in live mode.

Examples:

Stream historical BGP data:

config = BGPStreamConfig(
    start_time=datetime.datetime(2010, 9, 1, 0, 0),
    end_time=datetime.datetime(2010, 9, 1, 2, 0),
    collectors=["route-views.wide"],
)
stream = BGPStream.from_config(config)
for elem in stream:
    print(elem)

Direct instantiation with filters:

stream = BGPStream(
    collectors=["route-views.wide"],
    data_type=["updates"],
    ts_start=1283203200,
    ts_end=1283289600,
    filters=FilterOptions(origin_asn=64512),
    parser_name="bgpkit",
)
for elem in stream:
    print(f"{elem.prefix}: {elem.fields['as-path']}")

Live streaming from RIS Live:

config = BGPStreamConfig(
    collectors=["rrc00"],
    data_types=["updates"],
)
stream = BGPStream.from_config(config)
for elem in stream:
    print(f"Live: {elem.type} {elem.prefix}")
Source code in src/pybgpflux/bgpstream.py
 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
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
class BGPStream:
    """Stream and process BGP messages from multiple collectors.

    BGPStream is a high-performance alternative to PyBGPStream that parses BGP
    MRT files using BGPKIT. It can stream both historical and live BGP data with
    support for advanced filtering, multiple parser backends, and memory-efficient
    lazy loading.

    Attributes:
        collectors (list[str]): List of collector names to fetch data from.
        data_type (list[Literal["ribs", "updates"]]): Data types to stream ("ribs" or "updates").
        ts_start (float | None): Start timestamp (Unix epoch). None for live mode.
        ts_end (float | None): End timestamp (Unix epoch). None for live mode.
        filters (FilterOptions): Filtering options for BGP elements.
        cache_dir (Directory | TemporaryDirectory): Cache directory for downloaded files.
        parser_name (str): Backend parser to use ("pybgpkit", "bgpkit", "bgpdump", "pybgpstream").
        max_concurrent_downloads (int): Maximum concurrent file downloads.
        ram_fetch (bool): Use RAM disk (/dev/shm, /Volumes/RAMDisk) if available.
        jitter_buffer_delay (float): Delay (seconds) for jitter buffer in live mode.

    Examples:
        Stream historical BGP data:

        ```python
        config = BGPStreamConfig(
            start_time=datetime.datetime(2010, 9, 1, 0, 0),
            end_time=datetime.datetime(2010, 9, 1, 2, 0),
            collectors=["route-views.wide"],
        )
        stream = BGPStream.from_config(config)
        for elem in stream:
            print(elem)
        ```

        Direct instantiation with filters:

        ```python
        stream = BGPStream(
            collectors=["route-views.wide"],
            data_type=["updates"],
            ts_start=1283203200,
            ts_end=1283289600,
            filters=FilterOptions(origin_asn=64512),
            parser_name="bgpkit",
        )
        for elem in stream:
            print(f"{elem.prefix}: {elem.fields['as-path']}")
        ```

        Live streaming from RIS Live:

        ```python
        config = BGPStreamConfig(
            collectors=["rrc00"],
            data_types=["updates"],
        )
        stream = BGPStream.from_config(config)
        for elem in stream:
            print(f"Live: {elem.type} {elem.prefix}")
        ```
    """

    def __init__(
        self,
        collectors: list[str],
        data_types: list[Literal["ribs", "updates"]],
        ts_start: datetime.datetime | None = None,
        ts_end: datetime.datetime | None = None,
        filters: FilterOptions | None = None,
        cache_dir: str | None = None,
        max_concurrent_downloads: int | None = 10,
        ram_fetch: bool | None = True,
        parser_name: str | None = "pybgpkit",
        remote_parse: bool | None = True,
        broker: str | None = "bgpkit",
        jitter_buffer_delay: float | None = 10.0,
    ):
        """Initialize a BGP stream.

        Args:
            collectors: List of collector names (e.g., ["route-views.wide", "rrc04"]).
            data_types: List of data types to stream ("update", "rib", or both).
            ts_start: Start timestamp (Unix epoch) for historical data. None for live mode.
            ts_end: End timestamp (Unix epoch) for historical data. None for live mode.
            filters: Optional FilterOptions to filter BGP elements. Defaults to no filtering.
            cache_dir: Directory to cache downloaded MRT files. If None, uses temporary directory.
            max_concurrent_downloads: Maximum concurrent downloads. Default is 10.
            ram_fetch: Use RAM disk for temporary files if available. Default is True.
            parser_name: Parser backend ("pybgpkit", "bgpkit", "bgpdump", "pybgpstream").
                Default is "pybgpkit" (no system dependencies).
            broker: Archive file broker to query. Default is "bgpkit".
            jitter_buffer_delay: Delay (seconds) for jitter buffer in live mode. Default is 10.0.

        Raises:
            ValueError: If parser_name is invalid.

        Note:
            For live mode, set both ts_start and ts_end to None.
            For historical data, both ts_start and ts_end must be provided.
        """
        # Stream config
        self.ts_start = ts_start
        self.ts_end = ts_end
        self.collectors = collectors
        self.data_types = data_types
        if not filters:
            filters = FilterOptions()
        self.filters = filters

        # Implementation config
        if max_concurrent_downloads:
            self.max_concurrent_downloads = max_concurrent_downloads
        else:
            self.max_concurrent_downloads = 10

        self.ram_fetch = ram_fetch
        self.cache_dir = cache_dir

        if not parser_name:
            self.parser_name = "pybgpkit"
        else:
            self.parser_name = parser_name

        if not remote_parse:
            self.remote_parse = True
        else:
            self.remote_parse = remote_parse

        if cache_dir:
            self.remote_parse = False

        match broker:
            case "bgpkit":
                self.broker = BGPKITBroker()
            case "bgpstream":
                self.broker = BGPStreamBroker(
                    url="https://broker.bgpstream.caida.org/v2"
                )
            case "bgpfinder":
                self.broker = BGPStreamBroker(
                    url="https://bgpfinder.inetintel.cc.gatech.edu"
                )
            case _:
                self.broker = BGPKITBroker()

        self.parser_cls: type[BGPParser] = name2parser[self.parser_name]

        # Live config
        self.jitter_buffer_delay = jitter_buffer_delay

    def _set_urls(self):
        """Set archive files URL with a broker and setup prefetch queues"""
        self.urls: RCUrlsWithQueues = {
            "ribs": defaultdict(lambda: ([], asyncio.Queue(maxsize=PREFETCH_SIZE))),
            "updates": defaultdict(lambda: ([], asyncio.Queue(maxsize=PREFETCH_SIZE))),
        }
        config = BGPStreamConfig(
            start_time=self.ts_start,
            end_time=self.ts_end,
            collectors=self.collectors,
            data_types=self.data_types,
        )
        items = self.broker.query(config)
        items.sort(key=lambda item: dt_from_filepath(item.url))

        for item in items:
            self.urls[item.data_type][item.collector_id][0].append(item.url)

    def __iter__(self):
        if self.ts_start is None and self.ts_end is None:
            return self._iter_live()
        return self._iter_archive()

    @staticmethod
    def _run_loop(loop: asyncio.AbstractEventLoop, ready: threading.Event):
        """Background thread: run the event loop forever until stopped."""
        asyncio.set_event_loop(loop)
        loop.call_soon(ready.set)
        loop.run_forever()

    def _iter_archive(self) -> Iterator[BGPElement]:
        """__iter__ for data types [ribs, updates] or [updates]"""
        assert self.ts_start
        assert self.ts_end

        if cache_dir := self.cache_dir:
            is_caching = True
            cache_dir = Directory(self.cache_dir)
        else:
            # Note that if the parser supports remote parsing, cache_dir will not be populated
            is_caching = False
            if self.ram_fetch:
                cache_dir = TemporaryDirectory(dir=get_shared_memory())
            else:
                cache_dir = TemporaryDirectory()

        with cache_dir:
            self._set_urls()

            loop = asyncio.new_event_loop()

            # Single background thread runs the event loop
            ready = threading.Event()
            bg_thread = threading.Thread(
                target=self._run_loop, args=(loop, ready), daemon=True
            )
            bg_thread.start()
            ready.wait()

            is_remote_parsing = (
                not is_caching
                and self.parser_cls.supports_remote_parsing
                and self.remote_parse
            )

            # Kick off all download tasks in the background thread
            asyncio.run_coroutine_threadsafe(
                safe_download_all(
                    self.urls,
                    cache_dir.name,
                    self.max_concurrent_downloads,
                    is_remote_parsing,
                ),
                loop,
            )

            streams = [
                RCStream(
                    self.parser_cls,
                    data_type,
                    collector,
                    self.filters,
                    is_caching,
                    is_remote_parsing,
                    async_q,
                    loop,
                )
                for data_type in self.urls
                for collector, (_, async_q) in self.urls[data_type].items()
            ]

            try:
                ts_start, ts_end = self.ts_start.timestamp(), self.ts_end.timestamp()
                for bgpelem in merge(*streams, key=attrgetter("time")):
                    if ts_start <= bgpelem.time <= ts_end:
                        yield bgpelem
            finally:
                # Cancel in-flight download tasks first so workers blocked on a
                # full prefetch queue are released, then stop and close the loop.
                try:
                    asyncio.run_coroutine_threadsafe(cancel_all_tasks(), loop).result(
                        timeout=5.0
                    )
                except Exception as e:
                    logger.debug(f"Task cancellation during teardown failed: {e}")
                loop.call_soon_threadsafe(loop.stop)
                bg_thread.join(timeout=5.0)
                loop.close()

    def _iter_live(self) -> Iterator[BGPElement]:
        ris_collectors = [
            collector for collector in self.collectors if collector[:3] == "rrc"
        ]

        stream = RISLiveStream(collectors=ris_collectors, filters=self.filters)

        if self.jitter_buffer_delay is not None and self.jitter_buffer_delay > 0:
            stream = jitter_buffer_stream(stream, buffer_delay=self.jitter_buffer_delay)

        for elem in stream:
            yield elem

    @classmethod
    def from_config(cls, config: BGPStreamConfig | LiveStreamConfig) -> "BGPStream":
        """Create a BGPStream from a configuration object.

        Factory method to create a stream from various configuration types,
        automatically handling conversions and parameter mappings.

        Args:
            config: Configuration object, one of:
                - BGPStreamConfig: Unified configuration with query and implementation options.
                - LiveStreamConfig: Configuration for live RIS Live streaming.

        Returns:
            BGPStream: Initialized stream ready for iteration.

        Examples:
            ```python
            from pybgpflux import BGPStreamConfig, BGPStream
            import datetime

            config = BGPStreamConfig(
                start_time=datetime.datetime(2010, 9, 1, 0, 0),
                end_time=datetime.datetime(2010, 9, 1, 2, 0),
                collectors=["route-views.wide"],
            )
            stream = BGPStream.from_config(config)
            for elem in stream:
                print(elem)
            ```
        """
        match config:
            case BGPStreamConfig():
                if not config.is_live():
                    return cls(
                        ts_start=config.start_time,  # type: ignore
                        ts_end=config.end_time,  # type: ignore
                        collectors=config.collectors,
                        data_types=config.data_types,
                        filters=config.filters if config.filters else FilterOptions(),
                        cache_dir=str(config.cache_dir) if config.cache_dir else None,
                        max_concurrent_downloads=config.max_concurrent_downloads
                        if config.max_concurrent_downloads
                        else 10,
                        ram_fetch=config.ram_fetch if config.ram_fetch else None,
                        parser_name=config.parser if config.parser else "pybgpkit",
                        remote_parse=config.remote_parse
                        if config.remote_parse
                        else True,
                    )
                else:
                    return cls(
                        collectors=config.collectors,
                        data_types=["updates"],
                        filters=config.filters if config.filters else FilterOptions(),
                        jitter_buffer_delay=10,
                    )
            case LiveStreamConfig():
                return cls(
                    collectors=config.collectors,
                    data_types=["updates"],
                    filters=config.filters if config.filters else FilterOptions(),
                    jitter_buffer_delay=config.jitter_buffer_delay,
                )

            case _:
                raise ValueError("Unsupported config type")

__init__(collectors, data_types, ts_start=None, ts_end=None, filters=None, cache_dir=None, max_concurrent_downloads=10, ram_fetch=True, parser_name='pybgpkit', remote_parse=True, broker='bgpkit', jitter_buffer_delay=10.0)

Initialize a BGP stream.

Parameters:

Name Type Description Default
collectors list[str]

List of collector names (e.g., ["route-views.wide", "rrc04"]).

required
data_types list[Literal['ribs', 'updates']]

List of data types to stream ("update", "rib", or both).

required
ts_start datetime | None

Start timestamp (Unix epoch) for historical data. None for live mode.

None
ts_end datetime | None

End timestamp (Unix epoch) for historical data. None for live mode.

None
filters FilterOptions | None

Optional FilterOptions to filter BGP elements. Defaults to no filtering.

None
cache_dir str | None

Directory to cache downloaded MRT files. If None, uses temporary directory.

None
max_concurrent_downloads int | None

Maximum concurrent downloads. Default is 10.

10
ram_fetch bool | None

Use RAM disk for temporary files if available. Default is True.

True
parser_name str | None

Parser backend ("pybgpkit", "bgpkit", "bgpdump", "pybgpstream"). Default is "pybgpkit" (no system dependencies).

'pybgpkit'
broker str | None

Archive file broker to query. Default is "bgpkit".

'bgpkit'
jitter_buffer_delay float | None

Delay (seconds) for jitter buffer in live mode. Default is 10.0.

10.0

Raises:

Type Description
ValueError

If parser_name is invalid.

Note

For live mode, set both ts_start and ts_end to None. For historical data, both ts_start and ts_end must be provided.

Source code in src/pybgpflux/bgpstream.py
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
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
def __init__(
    self,
    collectors: list[str],
    data_types: list[Literal["ribs", "updates"]],
    ts_start: datetime.datetime | None = None,
    ts_end: datetime.datetime | None = None,
    filters: FilterOptions | None = None,
    cache_dir: str | None = None,
    max_concurrent_downloads: int | None = 10,
    ram_fetch: bool | None = True,
    parser_name: str | None = "pybgpkit",
    remote_parse: bool | None = True,
    broker: str | None = "bgpkit",
    jitter_buffer_delay: float | None = 10.0,
):
    """Initialize a BGP stream.

    Args:
        collectors: List of collector names (e.g., ["route-views.wide", "rrc04"]).
        data_types: List of data types to stream ("update", "rib", or both).
        ts_start: Start timestamp (Unix epoch) for historical data. None for live mode.
        ts_end: End timestamp (Unix epoch) for historical data. None for live mode.
        filters: Optional FilterOptions to filter BGP elements. Defaults to no filtering.
        cache_dir: Directory to cache downloaded MRT files. If None, uses temporary directory.
        max_concurrent_downloads: Maximum concurrent downloads. Default is 10.
        ram_fetch: Use RAM disk for temporary files if available. Default is True.
        parser_name: Parser backend ("pybgpkit", "bgpkit", "bgpdump", "pybgpstream").
            Default is "pybgpkit" (no system dependencies).
        broker: Archive file broker to query. Default is "bgpkit".
        jitter_buffer_delay: Delay (seconds) for jitter buffer in live mode. Default is 10.0.

    Raises:
        ValueError: If parser_name is invalid.

    Note:
        For live mode, set both ts_start and ts_end to None.
        For historical data, both ts_start and ts_end must be provided.
    """
    # Stream config
    self.ts_start = ts_start
    self.ts_end = ts_end
    self.collectors = collectors
    self.data_types = data_types
    if not filters:
        filters = FilterOptions()
    self.filters = filters

    # Implementation config
    if max_concurrent_downloads:
        self.max_concurrent_downloads = max_concurrent_downloads
    else:
        self.max_concurrent_downloads = 10

    self.ram_fetch = ram_fetch
    self.cache_dir = cache_dir

    if not parser_name:
        self.parser_name = "pybgpkit"
    else:
        self.parser_name = parser_name

    if not remote_parse:
        self.remote_parse = True
    else:
        self.remote_parse = remote_parse

    if cache_dir:
        self.remote_parse = False

    match broker:
        case "bgpkit":
            self.broker = BGPKITBroker()
        case "bgpstream":
            self.broker = BGPStreamBroker(
                url="https://broker.bgpstream.caida.org/v2"
            )
        case "bgpfinder":
            self.broker = BGPStreamBroker(
                url="https://bgpfinder.inetintel.cc.gatech.edu"
            )
        case _:
            self.broker = BGPKITBroker()

    self.parser_cls: type[BGPParser] = name2parser[self.parser_name]

    # Live config
    self.jitter_buffer_delay = jitter_buffer_delay

from_config(config) classmethod

Create a BGPStream from a configuration object.

Factory method to create a stream from various configuration types, automatically handling conversions and parameter mappings.

Parameters:

Name Type Description Default
config BGPStreamConfig | LiveStreamConfig

Configuration object, one of: - BGPStreamConfig: Unified configuration with query and implementation options. - LiveStreamConfig: Configuration for live RIS Live streaming.

required

Returns:

Name Type Description
BGPStream BGPStream

Initialized stream ready for iteration.

Examples:

from pybgpflux import BGPStreamConfig, BGPStream
import datetime

config = BGPStreamConfig(
    start_time=datetime.datetime(2010, 9, 1, 0, 0),
    end_time=datetime.datetime(2010, 9, 1, 2, 0),
    collectors=["route-views.wide"],
)
stream = BGPStream.from_config(config)
for elem in stream:
    print(elem)
Source code in src/pybgpflux/bgpstream.py
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
@classmethod
def from_config(cls, config: BGPStreamConfig | LiveStreamConfig) -> "BGPStream":
    """Create a BGPStream from a configuration object.

    Factory method to create a stream from various configuration types,
    automatically handling conversions and parameter mappings.

    Args:
        config: Configuration object, one of:
            - BGPStreamConfig: Unified configuration with query and implementation options.
            - LiveStreamConfig: Configuration for live RIS Live streaming.

    Returns:
        BGPStream: Initialized stream ready for iteration.

    Examples:
        ```python
        from pybgpflux import BGPStreamConfig, BGPStream
        import datetime

        config = BGPStreamConfig(
            start_time=datetime.datetime(2010, 9, 1, 0, 0),
            end_time=datetime.datetime(2010, 9, 1, 2, 0),
            collectors=["route-views.wide"],
        )
        stream = BGPStream.from_config(config)
        for elem in stream:
            print(elem)
        ```
    """
    match config:
        case BGPStreamConfig():
            if not config.is_live():
                return cls(
                    ts_start=config.start_time,  # type: ignore
                    ts_end=config.end_time,  # type: ignore
                    collectors=config.collectors,
                    data_types=config.data_types,
                    filters=config.filters if config.filters else FilterOptions(),
                    cache_dir=str(config.cache_dir) if config.cache_dir else None,
                    max_concurrent_downloads=config.max_concurrent_downloads
                    if config.max_concurrent_downloads
                    else 10,
                    ram_fetch=config.ram_fetch if config.ram_fetch else None,
                    parser_name=config.parser if config.parser else "pybgpkit",
                    remote_parse=config.remote_parse
                    if config.remote_parse
                    else True,
                )
            else:
                return cls(
                    collectors=config.collectors,
                    data_types=["updates"],
                    filters=config.filters if config.filters else FilterOptions(),
                    jitter_buffer_delay=10,
                )
        case LiveStreamConfig():
            return cls(
                collectors=config.collectors,
                data_types=["updates"],
                filters=config.filters if config.filters else FilterOptions(),
                jitter_buffer_delay=config.jitter_buffer_delay,
            )

        case _:
            raise ValueError("Unsupported config type")