Skip to content

Commit 2e1b336

Browse files
committed
feat(stream): handles over TA-Lib C 0.8.1's streaming API (#757)
talib.stream.SMA(close) returns a handle, not a value: .value is the value at the last history bar, .update(bar) is O(1) and returns that bar's, .peek(bar) evaluates a forming bar without committing, .copy() forks it, and .open_and_fill() returns the handle plus the Function API's series in one pass. Multi-output functions answer with a named tuple; .out_range and .advance() carry the C contract's bar range. BREAKING: talib.stream_X and the old value-returning talib.stream.X are gone, along with the stream_* stubs in _ta_lib.pyi -- talib/stream.pyi types the handles instead. Migrating is stream.X(...) -> stream.X(...).value, and the compiler will not find the sites: `if stream.CDLDOJI(o, h, l, c):` used to test the pattern and now tests a handle, which is always true. Adds talib.InsufficientHistory, raised when an open gets too little history. talib/_stream.pxi and talib/stream.pyi are generated by tools/generate_stream.py (--stub for the second).
1 parent ee0137c commit 2e1b336

19 files changed

Lines changed: 277002 additions & 59835 deletions

‎CHANGELOG‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,24 @@
33

44
- [NEW]: Support TA-Lib C 0.8.1, which is now the minimum required version.
55

6+
- [CHANGE]: ``talib.stream`` is now the real streaming API of TA-Lib C 0.8.1:
7+
``stream.SMA(close)`` returns a handle, not a value. ``handle.value`` is the
8+
value at the last history bar, ``handle.update(bar)`` costs O(1) and returns
9+
that bar's value, ``handle.peek(bar)`` evaluates a forming bar without
10+
committing it, and ``handle.copy()`` forks it. ``stream.SMA.open_and_fill()``
11+
returns the handle and the Function API's series in one pass. A multi-output
12+
function answers with a named tuple. The old last-value functions --
13+
``talib.stream.SMA``, ``talib.stream_SMA``, and their ``_ta_lib.pyi`` stubs --
14+
are gone; ``talib/stream.pyi`` types the handles instead.
15+
16+
Migrating is ``stream.X(...)`` -> ``stream.X(...).value``, and the compiler
17+
cannot find the sites for you: ``if stream.CDLDOJI(o, h, l, c):`` used to test
18+
the pattern and now tests a handle, which is always true.
19+
20+
- [NEW]: ``talib.InsufficientHistory``, raised when a stream is opened with
21+
too little history. It is the library's one recoverable error, so it is
22+
catchable on its own rather than as a bare ``Exception``.
23+
624
- [NEW]: The 40 functions TA-Lib C added since 0.7.1: AC, ADR, AO, CMF, CMOU,
725
COPPOCK, CUMSUM, CVI, DONCHIAN, DPO, EFI, ER, ERI, FOSC, FRACTAL, HA, HMA,
826
KC, KDJ, MARKETFI, MASSI, NVI, PERCENTILE, PERCENTRANK, PVI, PVO, PVT,

‎DEVELOPMENT‎

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -38,11 +38,15 @@ talib/_ta_lib.pyx
3838
need to use in the above pyx files.
3939

4040
talib/_stream.pxi
41-
This file contains code for interfacing a "streaming" interface to TA-Lib.
41+
This file is generated automatically by tools/generate_stream.py: one handle
42+
class per indicator, over TA-Lib C's streaming API.
43+
44+
talib/stream.pyi
45+
Type stubs for those handles, generated by tools/generate_stream.py --stub.
4246

4347
tools/generate_func.py,generate_stream.py
44-
Scripts that generate and print _func.pxi or _stream.pxi to stdout. Gets information
45-
about all functions from the C headers of the installed TA-Lib.
48+
Scripts that generate and print _func.pxi, _stream.pxi or stream.pyi to stdout.
49+
Gets information about all functions from the C headers of the installed TA-Lib.
4650

4751
If you are interested in developing new indicator functions or whatnot on
4852
the underlying TA-Lib, you must install TA-Lib from git.

‎MANIFEST.in‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,4 +5,5 @@ include talib/*.c
55
include talib/*.pyx
66
include talib/*.pxd
77
include talib/*.pxi
8+
include talib/*.pyi
89
include tests/*.py

‎Makefile‎

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,13 @@ talib/_func.pxi: tools/generate_func.py
1212
talib/_stream.pxi: tools/generate_stream.py
1313
python3 tools/generate_stream.py > talib/_stream.pxi
1414

15-
generate: talib/_func.pxi talib/_stream.pxi
15+
talib/stream.pyi: tools/generate_stream.py
16+
python3 tools/generate_stream.py --stub > talib/stream.pyi
17+
18+
talib/abstract.pyi: tools/generate_abstract_stub.py
19+
python3 tools/generate_abstract_stub.py > talib/abstract.pyi
20+
21+
generate: talib/_func.pxi talib/_stream.pxi talib/stream.pyi talib/abstract.pyi
1622

1723
cython:
1824
cython talib/_ta_lib.pyx

‎README.md‎

Lines changed: 66 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -532,27 +532,82 @@ slowk, slowd = STOCH(inputs, 5, 3, 0, 3, 0, prices=['high', 'low', 'open'])
532532
533533
## Streaming API
534534
535-
An experimental Streaming API was added that allows users to compute the latest
536-
value of an indicator. This can be faster than using the Function API, for
537-
example in an application that receives streaming data, and wants to know just
538-
the most recent updated indicator value.
535+
The Streaming API keeps a handle per indicator instead of recomputing from the
536+
whole array. Opening one costs a pass over the history; every bar after that is
537+
O(1), and each value it produces is identical to the one the Function API
538+
reports for that bar.
539539
540540
```python
541541
import talib
542542
from talib import stream
543543
544-
close = np.random.random(100)
545-
546-
# the Function API
544+
# the Function API: the whole series, from the whole array
547545
output = talib.SMA(close)
548546
549-
# the Streaming API
550-
latest = stream.SMA(close)
547+
# the Streaming API: a handle, positioned at the end of the history
548+
s = stream.SMA(close)
549+
assert s.value == output[-1]
551550
552-
# the latest value is the same as the last output value
553-
assert (output[-1] - latest) < 0.00001
551+
for price in feed:
552+
latest = s.update(price) # one closed bar in, its value out
553+
554+
s.peek(forming) # what update would return, committing nothing
555+
fork = s.copy() # an independent handle at the same bar
556+
```
557+
558+
`stream.SMA` takes exactly the arguments `talib.SMA` takes. A single-output
559+
function answers with a `float` (an `int` where the Function API returns an
560+
integer array); a multi-output one with a named tuple that still unpacks like
561+
the Function API's tuple:
562+
563+
```python
564+
m = stream.MACD(close)
565+
macd, macdsignal, macdhist = m.update(price)
566+
m.value.macdhist
554567
```
555568
569+
Opening needs at least `lookback + 1` bars, which `abstract` knows, and a little
570+
more where a function's seeding does -- so rather than computing the number,
571+
treat a short history as "not yet":
572+
573+
```python
574+
from talib import abstract
575+
576+
need = abstract.Function('RSI', timeperiod=14).lookback + 1 # 15, usually enough
577+
578+
try:
579+
s = stream.RSI(history, timeperiod=14)
580+
except talib.InsufficientHistory:
581+
... # collect more bars
582+
```
583+
584+
Leading bars that are NaN in any input are not history. They are skipped, as the
585+
Function API skips them, and do not count toward the warm-up. A NaN or an
586+
infinity anywhere else in the history is undefined behaviour in TA-Lib C, and
587+
for a few window functions a handle and the Function API do then disagree.
588+
589+
A bar that is not finite is likewise rejected: `update` raises and the handle is
590+
left exactly as it was, neither its value nor its range moved. For a bar you
591+
mean to skip rather than re-feed, say so with `advance()`, or two handles on one
592+
feed drift a bar apart.
593+
594+
If you want the series over the history as well, one pass gives both:
595+
596+
```python
597+
s, rsi = stream.RSI.open_and_fill(history, timeperiod=14) # rsi == talib.RSI(history)
598+
```
599+
600+
A handle also reports the range it has an output for, in the input series'
601+
coordinates, and can be told about a bar it was not fed:
602+
603+
```python
604+
s.out_range # (begidx, nbelement), as the Function API's output
605+
s.advance() # count a skipped bar: the range moves, the value holds
606+
```
607+
608+
A handle points into the TA-Lib C library, so it cannot be pickled or shared
609+
with another process. Keep the history and re-open instead.
610+
556611
## Supported Indicators and Functions 📋
557612
558613
We can show all the TA functions supported by TA-Lib, either as a `list` or

‎pyproject.toml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -40,4 +40,4 @@ requires-python = '>=3.9'
4040

4141
[tool.setuptools]
4242
packages = ["talib"]
43-
package-data = {"talib" = ["_ta_lib.pyi", "py.typed", "abstract.pyi"]}
43+
package-data = {"talib" = ["_ta_lib.pyi", "py.typed", "abstract.pyi", "stream.pyi"]}

‎talib/__init__.py‎

Lines changed: 3 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -79,12 +79,6 @@ def wrapper(*args, **kwds):
7979

8080
result = func(*_args, **_kwds)
8181

82-
# check to see if we got a streaming result
83-
first_result = result[0] if isinstance(result, tuple) else result
84-
is_streaming_fn_result = not hasattr(first_result, '__len__')
85-
if is_streaming_fn_result:
86-
return result
87-
8882
# Series was passed in, Series gets out
8983
if use_pl:
9084
if isinstance(result, tuple):
@@ -116,6 +110,7 @@ def wrapper(*args, **kwds):
116110
_ta_get_unstable_period as get_unstable_period,
117111
_ta_set_compatibility as set_compatibility,
118112
_ta_get_compatibility as get_compatibility,
113+
InsufficientHistory,
119114
__TA_FUNCTION_NAMES__
120115
)
121116
except ImportError as error:
@@ -141,12 +136,7 @@ def wrapper(*args, **kwds):
141136
setattr(func, func_name, wrapped_func)
142137
globals()[func_name] = wrapped_func
143138

144-
stream_func_names = ['stream_%s' % fname for fname in __TA_FUNCTION_NAMES__]
145-
stream = __import__("stream", globals(), locals(), stream_func_names, level=1)
146-
for func_name, stream_func_name in zip(__TA_FUNCTION_NAMES__, stream_func_names):
147-
wrapped_func = _wrapper(getattr(stream, func_name))
148-
setattr(stream, func_name, wrapped_func)
149-
globals()[stream_func_name] = wrapped_func
139+
from . import stream
150140

151141
__version__ = '0.7.1'
152142

@@ -400,4 +390,4 @@ def get_function_groups():
400390
"""
401391
return __function_groups__.copy()
402392

403-
__all__ = ['get_functions', 'get_function_groups'] + __TA_FUNCTION_NAMES__ + ["stream_%s" % name for name in __TA_FUNCTION_NAMES__]
393+
__all__ = ['get_functions', 'get_function_groups', 'InsufficientHistory', 'stream'] + __TA_FUNCTION_NAMES__

‎talib/_common.pxi‎

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,13 @@ from _ta_lib cimport TA_RetCode, TA_FuncUnstId
55

66
__ta_version__ = lib.TA_GetVersionString()
77

8+
9+
class InsufficientHistory(Exception):
10+
"""Not enough history to open a stream (TA_INSUFFICIENT_HISTORY).
11+
12+
Recoverable: collect more bars and try again."""
13+
14+
815
cpdef _ta_check_success(str function_name, TA_RetCode ret_code):
916
if ret_code == 0:
1017
return True
@@ -48,7 +55,8 @@ cpdef _ta_check_success(str function_name, TA_RetCode ret_code):
4855
description = 'Unknown Error (TA_UNKNOWN_ERR)'
4956
else:
5057
description = 'Unknown Error'
51-
raise Exception('%s function failed with error code %s: %s' % (
58+
error = InsufficientHistory if ret_code == 17 else Exception
59+
raise error('%s function failed with error code %s: %s' % (
5260
function_name, ret_code, description))
5361

5462
def _ta_initialize():

0 commit comments

Comments
 (0)