> ## Documentation Index
> Fetch the complete documentation index at: https://mintlify-poc.mintlify.site/llms.txt
> Use this file to discover all available pages before exploring further.

# Ingest real-time financial data

> Set up a data pipeline to get data from different financial APIs

export const CLOUD_LONG = 'Tiger Cloud';

export const CAGG = 'continuous aggregate';

export const SELF_LONG = 'self-hosted TimescaleDB';

export const SERVICE_SHORT = 'service';

export const COMPANY = 'Tiger Data ';

export const SERVICE_LONG = 'Tiger Cloud service';

export const CONSOLE = 'Tiger Console';

export const PG = 'Postgres';

export const HYPERTABLE_CAP = 'Hypertable';

export const HYPERCORE = 'hypercore';

export const ROWSTORE = 'rowstore';

export const HYPERCORE_CAP = 'Hypercore';

export const CHUNK = 'chunk';

export const HYPERTABLE = 'hypertable';

export const COLUMNSTORE = 'columnstore';

export const TIMESCALE_DB = 'TimescaleDB';

The financial industry is extremely data-heavy and relies on real-time and historical data for decision-making, risk assessment, fraud detection, and market analysis. {COMPANY} simplifies management of these large volumes of data, while also providing you with meaningful analytical insights and optimizing storage costs.

This tutorial shows you how to ingest real-time time-series data into
{TIMESCALE_DB} using a websocket connection. The tutorial sets up a data pipeline
to ingest real-time data from our data partner, [Twelve Data][twelve-data].
Twelve Data provides a number of different financial APIs, including stock,
cryptocurrencies, foreign exchanges, and ETFs. It also supports websocket
connections in case you want to update your database frequently. With
websockets, you need to connect to the server, subscribe to symbols, and you can
start receiving data in real-time during market hours.

When you complete this tutorial, you'll have a data pipeline set
up that ingests real-time financial data into your {CLOUD_LONG}.

This tutorial uses Python and the API
[wrapper library][twelve-wrapper] provided by Twelve Data.

This tutorial covers:

1. **Set up your dataset**: connect to the Twelve Data websocket server, create {HYPERTABLE}s, and ingest real-time cryptocurrency data.
2. **Query your data**: create {CAGG}s to aggregate OHLCV data, query the aggregated data, and visualize the data in Grafana.

## Prerequisites

To follow the steps on this page:

* Create a target [{SERVICE_LONG}][create-service] with Real-time analytics enabled.<p />

  You need [your connection details][connection-info]. This procedure also
  works for [{SELF_LONG}][enable-timescaledb].

[create-service]: /deploy-and-operate/tiger-cloud/get-started/create-services

[enable-timescaledb]: /deploy-and-operate/self-hosted/install-and-update/install-self-hosted

[connection-info]: /integrations/find-connection-details

* Install and run [self-managed Grafana][grafana-self-managed], or sign up for [Grafana Cloud][grafana-cloud].
* Install Python 3
* Sign up for [Twelve Data][twelve-signup]. The free tier is perfect for
  this tutorial.
* Made a note of your Twelve Data [API key][api-key].

## About OHLCV data and candlestick charts

The financial sector regularly uses [candlestick charts][charts] to visualize
the price change of an asset. Each candlestick represents a time period, such as
one minute or one hour, and shows how the asset's price changed during that time.

Candlestick charts are generated from the open, high, low, close, and volume
data for each financial asset during the time period. This is often abbreviated
as OHLCV:

* Open: opening price
* High: highest price
* Low: lowest price
* Close: closing price
* Volume: volume of transactions

[charts]: https://www.investopedia.com/terms/c/candlestick.asp

![candlestick][candlestick]

{TIMESCALE_DB} is well suited to storing and analyzing financial candlestick data,
and many {COMPANY} community members use it for exactly this purpose.

## Ingest data into a service

This tutorial uses a dataset that contains second-by-second cryptocurrency trade data,
in a {HYPERTABLE} named `crypto_ticks`. It also includes a separate table of
cryptocurrency symbols and names, in a regular {PG} table named `crypto_assets`.

### Connect to the websocket server

When you connect to the Twelve Data API through a websocket, you create a
persistent connection between your computer and the websocket server.
You set up a Python environment, and pass two arguments to create a
websocket object and establish the connection.

#### Set up a new Python environment

Create a new Python virtual environment for this project and activate it. All
the packages you need to complete for this tutorial are installed in this environment.

1. Create and activate a Python virtual environment:

   ```bash theme={"dark"}
   virtualenv env
   source env/bin/activate
   ```

2. Install the Twelve Data Python
   [wrapper library][twelve-wrapper]
   with websocket support. This library allows you to make requests to the
   API and maintain a stable websocket connection.

   ```bash theme={"dark"}
   pip install twelvedata websocket-client
   ```

3. Install [Psycopg2][psycopg2] so that you can connect the
   {TIMESCALE_DB} from your Python script:

   ```bash theme={"dark"}
   pip install psycopg2-binary
   ```

#### Create the websocket connection

A persistent connection between your computer and the websocket server is used
to receive data for as long as the connection is maintained. You need to pass
two arguments to create a websocket object and establish connection.

**Websocket arguments**

* `on_event`

  This argument needs to be a function that is invoked whenever there's a
  new data record is received from the websocket:

  ```python theme={"dark"}
  def on_event(event):
      print(event) # prints out the data record (dictionary)
  ```

  This is where you want to implement the ingestion logic so whenever
  there's new data available you insert it into the database.

* `symbols`

  This argument needs to be a list of stock ticker symbols (for example,
  `MSFT`) or crypto trading pairs (for example, `BTC/USD`). When using a
  websocket connection you always need to subscribe to the events you want to
  receive. You can do this by using the `symbols` argument or if your
  connection is already created you can also use the `subscribe()` function to
  get data for additional symbols.

**Connect to the websocket server**

1. Create a new Python file called `websocket_test.py` and connect to the
   Twelve Data servers using the `<YOUR_API_KEY>`:

   ```python theme={"dark"}
      import time
      from twelvedata import TDClient

       messages_history = []

       def on_event(event):
        print(event) # prints out the data record (dictionary)
        messages_history.append(event)

      td = TDClient(apikey="<YOUR_API_KEY>")
      ws = td.websocket(symbols=["BTC/USD", "ETH/USD"], on_event=on_event)
      ws.subscribe(['ETH/BTC', 'AAPL'])
      ws.connect()
      while True:
      print('messages received: ', len(messages_history))
      ws.heartbeat()
      time.sleep(10)
   ```

2. Run the Python script:

   ```bash theme={"dark"}
   python websocket_test.py
   ```

3. When you run the script, you receive a response from the server about the
   status of your connection:

   ```bash theme={"dark"}
   {'event': 'subscribe-status',
    'status': 'ok',
    'success': [
           {'symbol': 'BTC/USD', 'exchange': 'Coinbase Pro', 'mic_code': 'Coinbase Pro', 'country': '', 'type': 'Digital Currency'},
           {'symbol': 'ETH/USD', 'exchange': 'Huobi', 'mic_code': 'Huobi', 'country': '', 'type': 'Digital Currency'}
       ],
    'fails': None
   }
   ```

   When you have established a connection to the websocket server,
   wait a few seconds, and you can see data records, like this:

   ```bash theme={"dark"}
   {'event': 'price', 'symbol': 'BTC/USD', 'currency_base': 'Bitcoin', 'currency_quote': 'US Dollar', 'exchange': 'Coinbase Pro', 'type': 'Digital Currency', 'timestamp': 1652438893, 'price': 30361.2, 'bid': 30361.2, 'ask': 30361.2, 'day_volume': 49153}
   {'event': 'price', 'symbol': 'BTC/USD', 'currency_base': 'Bitcoin', 'currency_quote': 'US Dollar', 'exchange': 'Coinbase Pro', 'type': 'Digital Currency', 'timestamp': 1652438896, 'price': 30380.6, 'bid': 30380.6, 'ask': 30380.6, 'day_volume': 49157}
   {'event': 'heartbeat', 'status': 'ok'}
   {'event': 'price', 'symbol': 'ETH/USD', 'currency_base': 'Ethereum', 'currency_quote': 'US Dollar', 'exchange': 'Huobi', 'type': 'Digital Currency', 'timestamp': 1652438899, 'price': 2089.07, 'bid': 2089.02, 'ask': 2089.03, 'day_volume': 193818}
   {'event': 'price', 'symbol': 'BTC/USD', 'currency_base': 'Bitcoin', 'currency_quote': 'US Dollar', 'exchange': 'Coinbase Pro', 'type': 'Digital Currency', 'timestamp': 1652438900, 'price': 30346.0, 'bid': 30346.0, 'ask': 30346.0, 'day_volume': 49167}
   ```

   Each price event gives you multiple data points about the given trading pair
   such as the name of the exchange, and the current price. You can also
   occasionally see `heartbeat` events in the response; these events signal
   the health of the connection over time.
   At this point the websocket connection is working successfully to pass data.

## Optimize time-series data in a hypertable

{HYPERTABLE_CAP}s are {PG} tables in {TIMESCALE_DB} that automatically partition your time-series data by time. Time-series data represents the way a system, process, or behavior changes over time. {HYPERTABLE_CAP}s enable {TIMESCALE_DB} to work efficiently with time-series data. Each {HYPERTABLE} is made up of child tables called chunks. Each chunk is assigned a range of time, and only contains data from that range. When you run a query, {TIMESCALE_DB} identifies the correct chunk and runs the query on it, instead of going through the entire table.

[{HYPERCORE_CAP}][hypercore] is the hybrid row-columnar storage engine in {TIMESCALE_DB} used by {HYPERTABLE}s. Traditional
databases force a trade-off between fast inserts (row-based storage) and efficient analytics
(columnar storage). {HYPERCORE_CAP} eliminates this trade-off, allowing real-time analytics without sacrificing
transactional capabilities.

{HYPERCORE_CAP} dynamically stores data in the most efficient format for its lifecycle:

![Move from rowstore to columstore in hypercore][move-from-rowstore-to-columstore-in-hypercore]

* **Row-based storage for recent data**: the most recent chunk (and possibly more) is always stored in the {ROWSTORE},
  ensuring fast inserts, updates, and low-latency single record queries. Additionally, row-based storage is used as a
  writethrough for inserts and updates to columnar storage.
* **Columnar storage for analytical performance**: chunks are automatically compressed into the {COLUMNSTORE}, optimizing
  storage efficiency and accelerating analytical queries.

Unlike traditional columnar databases, {HYPERCORE} allows data to be inserted or modified at any stage, making it a
flexible solution for both high-ingest transactional workloads and real-time analytics—within a single database.

[hypercore]: /manage-data/capabilities/hypercore/understand-hypercore

[move-from-rowstore-to-columstore-in-hypercore]: https://assets.timescale.com/docs/images/hypercore_intro.svg

Because {TIMESCALE_DB} is 100% {PG}, you can use all the standard {PG} tables, indexes, stored procedures, and other objects alongside your {HYPERTABLE}s. This makes creating and working with {HYPERTABLE}s similar to standard {PG}.

1. **Connect to your {SERVICE_LONG}**

   In [{CONSOLE}][services-portal] open an [SQL editor][in-console-editors]. You can also connect to your service using [psql][psql].

2. **Create a {HYPERTABLE} to store the real-time cryptocurrency data**

   Create a [{HYPERTABLE}][hypertables-section] for your time-series data using [CREATE TABLE][hypertable-create-table].
   For [efficient queries][secondary-indexes] on data in the {COLUMNSTORE}, remember to `segmentby` the column you will
   use most often to filter your data:

   ```sql theme={"dark"}
   CREATE TABLE crypto_ticks (
       "time" TIMESTAMPTZ,
       symbol TEXT,
       price DOUBLE PRECISION,
       day_volume NUMERIC
   ) WITH (
      tsdb.hypertable,
      tsdb.segmentby='symbol',
      tsdb.orderby='time DESC'
   );
   ```

   When you create a {HYPERTABLE} using [CREATE TABLE ... WITH ...][hypertable-create-table], the default partitioning
   column is automatically the first column with a timestamp data type. Also, {TIMESCALE_DB} creates a
   [columnstore policy][add_columnstore_policy] that automatically converts your data to the {COLUMNSTORE}, after an
   interval equal to the value of the [chunk\_interval][create_table_arguments], defined through `compress_after` in the
   policy. This columnar format enables fast scanning and
   aggregation, optimizing performance for analytical workloads while also saving significant storage space. In the
   {COLUMNSTORE} conversion, {HYPERTABLE} {CHUNK}s are compressed by up to 98%, and organized for efficient, large-scale queries.

   You can customize this policy later using [alter\_job][alter_job_samples]. However, to change `after` or
   `created_before`, the compression settings, or the {HYPERTABLE} the policy is acting on, you must
   [remove the columnstore policy][remove_columnstore_policy] and [add a new one][add_columnstore_policy].

   You can also manually [convert {CHUNK}s][convert_to_columnstore] in a {HYPERTABLE} to the {COLUMNSTORE}.

   [add_columnstore_policy]: /api-reference/timescaledb/hypercore/add_columnstore_policy

   [alter_job_samples]: /api-reference/timescaledb/jobs-automation/alter_job#samples

   [convert_to_columnstore]: /api-reference/timescaledb/hypercore/convert_to_columnstore

   [create_table_arguments]: /api-reference/timescaledb/hypertables/create_table#arguments

   [hypertable-create-table]: /api-reference/timescaledb/hypertables/create_table

   [remove_columnstore_policy]: /api-reference/timescaledb/hypercore/remove_columnstore_policy

## Create a standard Postgres table for relational data

When you have relational data that enhances your time-series data, store that data in
standard {PG} relational tables.

1. **Add a table to store the asset symbol and name in a relational table**

   ```sql theme={"dark"}
   CREATE TABLE crypto_assets (
       symbol TEXT UNIQUE,
       "name" TEXT
   );
   ```

You now have two tables within your {SERVICE_LONG}. A hypertable named `crypto_ticks`, and a normal
{PG} table named `crypto_assets`.

[hypertable-create-table]: /api-reference/timescaledb/hypertables/create_table

[hypercore]: /manage-data/data-management/hypercore/understand-hypercore

[hypertables-section]: /manage-data/data-management/hypertables/understand-hypertables

[in-console-editors]: /deploy-and-operate/tiger-cloud/get-started/run-queries-from-console

[psql]: /integrations/psql

[secondary-indexes]: /use-timescale/hypercore/secondary-indexes

[services-portal]: https://console.cloud.timescale.com/dashboard/services

When you ingest data into a transactional database like {TIMESCALE_DB}, it is more
efficient to insert data in batches rather than inserting data row-by-row. Using
one transaction to insert multiple rows can significantly increase the overall
ingest capacity and speed of your {SERVICE_LONG}.

## Batching in memory

A common practice to implement batching is to store new records in memory
first, then after the batch reaches a certain size, insert all the records
from memory into the database in one transaction. The perfect batch size isn't
universal, but you can experiment with different batch sizes
(for example, 100, 1000, 10000, and so on) and see which one fits your use case better.
Using batching is a fairly common pattern when ingesting data into {TIMESCALE_DB}
from Kafka, Kinesis, or websocket connections.

To ingest the data into your {SERVICE_LONG}, you need to implement the
`on_event` function.

After the websocket connection is set up, you can use the `on_event` function
to ingest data into the database. This is a data pipeline that ingests real-time
financial data into your {SERVICE_LONG}.

You can implement a batching solution in Python with Psycopg2.
You can implement the ingestion logic within the `on_event` function that
you can then pass over to the websocket object.

This function needs to:

1. Check if the item is a data item, and not websocket metadata.
2. Adjust the data so that it fits the database schema, including the data
   types, and order of columns.
3. Add it to the in-memory batch, which is a list in Python.
4. If the batch reaches a certain size, insert the data, and reset or empty the list.

## Ingest data in real-time

1. Update the Python script that prints out the current batch size, so you can
   follow when data gets ingested from memory into your database. Use
   the `<HOST>`, `<PASSWORD>`, and `<PORT>` details for the {SERVICE_LONG}
   where you want to ingest the data and your API key from Twelve Data:

   ```python theme={"dark"}
   import time
   import psycopg2

   from twelvedata import TDClient
   from psycopg2.extras import execute_values
   from datetime import datetime

   class WebsocketPipeline():
       # name of the hypertable
       DB_TABLE = "crypto_ticks"

       # columns in the hypertable in the correct order
       DB_COLUMNS=["time", "symbol", "price", "day_volume"]

       # batch size used to insert data in batches
       MAX_BATCH_SIZE=100

       def __init__(self, conn):
           """Connect to the Twelve Data web socket server and stream
           data into the database.

           Args:
               conn: psycopg2 connection object
           """
           self.conn = conn
           self.current_batch = []
           self.insert_counter = 0

       def _insert_values(self, data):
           if self.conn is not None:
               cursor = self.conn.cursor()
               sql = f"""
               INSERT INTO {self.DB_TABLE} ({','.join(self.DB_COLUMNS)})
               VALUES %s;"""
               execute_values(cursor, sql, data)
               self.conn.commit()

       def _on_event(self, event):
           """This function gets called whenever there's a new data record coming
           back from the server.

           Args:
               event (dict): data record
           """
           if event["event"] == "price":
               # data record
               timestamp = datetime.utcfromtimestamp(event["timestamp"])
               data = (timestamp, event["symbol"], event["price"], event.get("day_volume"))

               # add new data record to batch
               self.current_batch.append(data)
               print(f"Current batch size: {len(self.current_batch)}")

               # ingest data if max batch size is reached then reset the batch
               if len(self.current_batch) == self.MAX_BATCH_SIZE:
                   self._insert_values(self.current_batch)
                   self.insert_counter += 1
                   print(f"Batch insert #{self.insert_counter}")
                   self.current_batch = []
           def start(self, symbols):
               """Connect to the web socket server and start streaming real-time data
               into the database.

               Args:
                   symbols (list of symbols): List of stock/crypto symbols
               """
               td = TDClient(apikey="<YOUR_API_KEY")
               ws = td.websocket(on_event=self._on_event)
               ws.subscribe(symbols)
               ws.connect()
               while True:
                  ws.heartbeat()
                  time.sleep(10)
       onn = psycopg2.connect(database="tsdb",
                           host="<HOST>",
                           user="tsdbadmin",
                           password="<PASSWORD>",
                           port="<PORT>")

       symbols = ["BTC/USD", "ETH/USD", "MSFT", "AAPL"]
       websocket = WebsocketPipeline(conn)
       websocket.start(symbols=symbols)
   ```

2. Run the script:

   ```bash theme={"dark"}
   python websocket_test.py
   ```

You can even create separate Python scripts to start multiple websocket
connections for different types of symbols, for example, one for stock, and
another one for cryptocurrency prices.

### Troubleshooting

If you see an error message similar to this:

```bash theme={"dark"}
2022-05-13 18:51:41,976 - ws-twelvedata - ERROR - TDWebSocket ERROR: Handshake status 200 OK
```

Then check that you use a proper API key received from Twelve Data.

## Query the data

To look at OHLCV values, the most effective way is to create a {CAGG}. You can create a {CAGG} to aggregate data
for each day, then set the aggregate to refresh every day, and aggregate
the last two days' worth of data.

### Creating a continuous aggregate

1. Connect to the {SERVICE_LONG} `tsdb` that contains the Twelve Data
   cryptocurrency dataset.

2. At the psql prompt, create the {CAGG} to aggregate data every
   day:

   ```sql theme={"dark"}
   CREATE MATERIALIZED VIEW one_day_candle
   WITH (timescaledb.continuous) AS
       SELECT
           time_bucket('1 day', time) AS bucket,
           symbol,
           FIRST(price, time) AS "open",
           MAX(price) AS high,
           MIN(price) AS low,
           LAST(price, time) AS "close",
           LAST(day_volume, time) AS day_volume
       FROM crypto_ticks
       GROUP BY bucket, symbol;
   ```

   When you create the {CAGG}, it refreshes by default.

3. Set a refresh policy to update the {CAGG} every day,
   if there is new data available in the {HYPERTABLE} for the last two days:

   ```sql theme={"dark"}
   SELECT add_continuous_aggregate_policy('one_day_candle',
       start_offset => INTERVAL '3 days',
       end_offset => INTERVAL '1 day',
       schedule_interval => INTERVAL '1 day');
   ```

### Query the continuous aggregate

When you have your {CAGG} set up, you can query it to get the
OHLCV values.

1. Connect to the {SERVICE_LONG} that contains the Twelve Data
   cryptocurrency dataset.

2. At the psql prompt, use this query to select all Bitcoin OHLCV data for the
   past 14 days, by time bucket:

   ```sql theme={"dark"}
   SELECT * FROM one_day_candle
   WHERE symbol = 'BTC/USD' AND bucket >= NOW() - INTERVAL '14 days'
   ORDER BY bucket;
   ```

   The result of the query looks like this:

   ```sql theme={"dark"}
            bucket         | symbol  |  open   |  high   |   low   |  close  | day_volume
   ------------------------+---------+---------+---------+---------+---------+------------
    2022-11-24 00:00:00+00 | BTC/USD |   16587 | 16781.2 | 16463.4 | 16597.4 |      21803
    2022-11-25 00:00:00+00 | BTC/USD | 16597.4 | 16610.1 | 16344.4 | 16503.1 |      20788
    2022-11-26 00:00:00+00 | BTC/USD | 16507.9 | 16685.5 | 16384.5 | 16450.6 |      12300
   ```

## Connect Grafana to Tiger Cloud

To visualize the results of your queries, enable Grafana to read the data in your {SERVICE_SHORT}:

1. **Log in to Grafana**

   In your browser, log in to either:

   * Self-hosted Grafana: at `http://localhost:3000/`. The default credentials are `admin`, `admin`.
   * Grafana Cloud: use the URL and credentials you set when you created your account.
2. **Add your {SERVICE_SHORT} as a data source**

   1. Open `Connections` > `Data sources`, then click `Add new data source`.

   2. Select `PostgreSQL` from the list.

   3. Configure the connection:
      * `Host URL`, `Database name`, `Username`, and `Password`

        Configure using your [connection details][connection-info]. `Host URL` is in the format `<host>:<port>`.
      * `TLS/SSL Mode`: select `require`.
      * `PostgreSQL options`: enable `TimescaleDB`.
      * Leave the default setting for all other fields.

   4. Click `Save & test`.

   Grafana checks that your details are set correctly.

[cloud-login]: https://console.cloud.timescale.com/

[connection-info]: /integrations/find-connection-details

[create-service]: /deploy-and-operate/tiger-cloud/get-started/create-services

[grafana-cloud]: https://grafana.com/get/

[grafana-self-managed]: https://grafana.com/get/?tab=self-managed

## Graph OHLCV data

When you have extracted the raw OHLCV data, you can use it to graph the result
in a candlestick chart, using Grafana.

1. In Grafana, from the `Dashboards` page, click `New` and select `New dashboard`.
2. Click `Add visualization`, then select the data source that connects to your {SERVICE_LONG} and the `Candlestick` visualization type in the top right.
3. In the `Queries` section, select `Code` and paste the query you used to get the OHLCV values:

   ```sql theme={"dark"}
   SELECT * FROM one_day_candle
   WHERE symbol = 'BTC/USD' AND bucket >= NOW() - INTERVAL '14 days'
   ORDER BY bucket;
   ```
4. Adjust elements of the table as required, and click `Apply` to save your
   graph to the dashboard.

   ![Creating a candlestick graph in Grafana using 1-day OHLCV tick data](https://assets.timescale.com/docs/images/Grafana_candlestick_1day.webp)

[api-key]: https://twelvedata.com/account/api-keys

[candlestick]: https://assets.timescale.com/docs/images/tutorials/intraday-stock-analysis/candlestick_fig.png

[grafana-cloud]: https://grafana.com/get/

[grafana-self-managed]: https://grafana.com/get/?tab=self-managed

[psycopg2]: https://www.psycopg.org/docs/

[twelve-data]: https://twelvedata.com

[twelve-signup]: https://twelvedata.com/pricing

[twelve-wrapper]: https://github.com/twelvedata/twelvedata-python
