/
cQube | Flow of Data

1. Ingestion Specification

Ingestion is a simple API that allows you to ingest your data as events. There are two steps for this.

  1. Define and Event Grammar

1{ 2 "instrument_details": { 3 "type": "COUNTER", 4 "key": "count" 5 }, 6 "name": "attendance_by_school_grade_for_a_day", 7 "is_active": true, 8 "event_schema": { 9 "$schema": "https://json-schema.org/draft/2019-09/schema", 10 "$id": "http://example.com/example.json", 11 "type": "object", 12 "default": {}, 13 "title": "attendance_by_school_grade_for_a_day", 14 "required": ["date", "grade", "school_id", "count"], 15 "properties": { 16 "grade": { 17 "type": "integer", 18 "default": 0, 19 "title": "The grade Schema", 20 "examples": [1] 21 }, 22 "school_id": { 23 "type": "integer", 24 "default": 0, 25 "title": "The school_id Schema", 26 "examples": [901] 27 } 28 }, 29 "examples": [ 30 { 31 "date": "2019-01-01", 32 "grade": 1, 33 "school_id": 901, 34 "count": 23 35 } 36 ] 37 } 38}

Instruments: Currently only a COUNTER instrument is supported.

2. Publish an event that satisfies the grammar

1{ 2 "date": "2019-01-01", 3 "grade": 1, 4 "school_id": 901, 5 "count": 23 6}

If the data is not aggregated at a daily level, an event could also be shared at a week, month, year level as well by just defining a start_date and end_date. Not the start_date and end_date should cleanly fall in a week, month, year, last 7 days, last 15 days, or last 30 days.

1{ 2 "start_date": "2019-01-01", 3 "end_date": "2019-02-01", 4 "grade": 1, 5 "school_id": 901, 6 "count": 23 7}

if you cannot guarantee this, it would be better to use an alternative API to push an event as below.

1{ 2 "week": "2019-W21", // or "month": "2019-M11" or "year": "2022" 3 "grade": 1, 4 "school_id": 901, 5 "count": 23 6}

Keywords and preferences

The following keywords in an event will be treated as special keywords and don’t need to be part of the schema - date, start_date, end_date, week, month, and year and are preferred in the same order. They cannot be mappings to a dimension or dataset.

2. Transformer Specification

The transformer is defined as an actor and can be sync as well as async. This allows you to write any transformer in any language as long as it gets the job done. For the current implementation in cQube we create custom processors for every single transformer in Nifi. Events are transformed in bulk every few mins (A config needs to be exposed to for defining how long do we wait for bulk transformations).

1const transformer = { 2 name: "transformer_name", 3 event_grammar: "es11", //EventGrammar on which the transformer can act. 4 dataset_grammar: "ds23", //Dataset that this transformer which modify. 5 config: { 6 actor((callback, transform, transformBulk) => { 7 8 // manage side effects here 9 // Examples - logging, pushing to an externa event bus, etc. 10 callback('SOME_EVENT'); 11 12 // receive from parent 13 transform(async (event: Event) => { 14 // transform event to dataset; Has access to event and dataset grammars for validation 15 }); 16 17 // for speedups 18 transformBulk(async (event: Event[]) => { 19 // transform events to dataset; Has access to event and dataset grammars for validation 20 }); 21 22 // disposal 23 return () => { 24 /* do cleanup here */ 25 /* push this to dataset */ 26 }; 27 }), 28 } 29} 30

How this event gets processed later is defined by a pipe. A pipe connects then event to a source. An example of a pipe is shown below

1{ 2 "event": "es11", 3 "transformer": "tr33", 4 "dataset": "ds23" 5}

You can also use keep the transformer as null to just bypass everything and not transform data and directly insert it in a dataset.

Defining Dimensions

cQube supports arbitrary dimensions to be stored. There are only two categories of dimensions that are supported:

  1. Time-based dimensions

  2. Dynamic dimensions

For the Attendance Example, the following dimensions are viable dimensions

  1. School (Dynamic)

  2. District (Dynamic)

  3. Time (Time Based) - week, month, year, last 7 days, last 15 days, last 30 days.

School as a dimension

Below is the specifications on how they should be defined in cQube.

  1. Define the grammar for how the dimension data needs to be stored. Given that Schools and Districts are very similar, they could be combined as well. For the school dimension, sample data looks like the following -

1{ 2 "data": { 3 "school_id": 901, 4 "name": "A green door", 5 "type": "GSSS", 6 "enrollement_count": 345, 7 "district": "District", 8 "block": "Block", 9 "cluster": "Cluster", 10 "village": "Village" 11 } 12}

The actual schema also includes the metadata that is required by cQube to store this dimension in the database.

  • The indexes field is used to create indexes on the database. This creates a simple BTree index for the list of columns shared.

  • The primary_id field is used to create a primary key on the database.

The corresponding JSON Schema would need to be entered to accept the dimension as a valid input

1{ 2 "$schema": "<https://json-schema.org/draft/2019-09/schema",> 3 "$id": "<http://example.com/example.json",> 4 "type": "object", 5 "default": {}, 6 "title": "Root Schema", 7 "required": [ 8 "name", 9 "type", 10 "storage", 11 "data" 12 ], 13 "properties": { 14 "name": { 15 "type": "string", 16 "default": "", 17 "title": "The name Schema", 18 "examples": [ 19 "Schools" 20 ] 21 }, 22 "type": { 23 "type": "string", 24 "default": "", 25 "title": "The type Schema", 26 "examples": [ 27 "dynamic" 28 ] 29 }, 30 "storage": { 31 "type": "object", 32 "default": {}, 33 "title": "The storage Schema", 34 "required": [ 35 "indexes", 36 "primaryId", 37 "retention", 38 "bucket_size" 39 ], 40 "properties": { 41 "indexes": { 42 "type": "array", 43 "default": [], 44 "title": "The indexes Schema", 45 "items": { 46 "type": "string", 47 "title": "A Schema", 48 "examples": [ 49 "name", 50 "type" 51 ] 52 }, 53 "examples": [ 54 ["name", 55 "type" 56 ] 57 ] 58 }, 59 "primaryId": { 60 "type": "string", 61 "default": "", 62 "title": "The primaryId Schema", 63 "examples": [ 64 "school_id" 65 ] 66 }, 67 "retention": { 68 "type": "null", 69 "default": null, 70 "title": "The retention Schema", 71 "examples": [ 72 null 73 ] 74 }, 75 "bucket_size": { 76 "type": "null", 77 "default": null, 78 "title": "The bucket_size Schema", 79 "examples": [ 80 null 81 ] 82 } 83 }, 84 "examples": [{ 85 "indexes": [ 86 "name", 87 "type" 88 ], 89 "primaryId": "school_id", 90 "retention": null, 91 "bucket_size": null 92 }] 93 }, 94 "data": { 95 "type": "object", 96 "default": {}, 97 "title": "The data Schema", 98 "required": [ 99 "school_id", 100 "name", 101 "type", 102 "enrollement_count", 103 "district", 104 "block", 105 "cluster", 106 "village" 107 ], 108 "properties": { 109 "school_id": { 110 "type": "integer", 111 "default": 0, 112 "title": "The school_id Schema", 113 "examples": [ 114 901 115 ] 116 }, 117 "name": { 118 "type": "string", 119 "default": "", 120 "title": "The name Schema", 121 "examples": [ 122 "A green door" 123 ] 124 }, 125 "type": { 126 "type": "string", 127 "default": "", 128 "title": "The type Schema", 129 "examples": [ 130 "GSSS" 131 ] 132 }, 133 "enrollement_count": { 134 "type": "integer", 135 "default": 0, 136 "title": "The enrollement_count Schema", 137 "examples": [ 138 345 139 ] 140 }, 141 "district": { 142 "type": "string", 143 "default": "", 144 "title": "The district Schema", 145 "examples": [ 146 "District" 147 ] 148 }, 149 "block": { 150 "type": "string", 151 "default": "", 152 "title": "The block Schema", 153 "examples": [ 154 "Block" 155 ] 156 }, 157 "cluster": { 158 "type": "string", 159 "default": "", 160 "title": "The cluster Schema", 161 "examples": [ 162 "Cluster" 163 ] 164 }, 165 "village": { 166 "type": "string", 167 "default": "", 168 "title": "The village Schema", 169 "examples": [ 170 "Village" 171 ] 172 } 173 }, 174 "examples": [{ 175 "school_id": 901, 176 "name": "A green door", 177 "type": "GSSS", 178 "enrollement_count": 345, 179 "district": "District", 180 "block": "Block", 181 "cluster": "Cluster", 182 "village": "Village" 183 }] 184 } 185 }, 186 "examples": [{ 187 "name": "Schools", 188 "type": "dynamic", 189 "storage": { 190 "indexes": [ 191 "name", 192 "type" 193 ], 194 "primaryId": "school_id", 195 "retention": null, 196 "bucket_size": null 197 }, 198 "data": { 199 "school_id": 901, 200 "name": "A green door", 201 "type": "GSSS", 202 "enrollement_count": 345, 203 "district": "District", 204 "block": "Block", 205 "cluster": "Cluster", 206 "village": "Village" 207 } 208 }] 209}
  1. Once the dimension is defined, it can be added to cQube by using the dimension schema (definition grammar) API.

  2. The next step is to insert data in the dimensions which can be done through a POST request to the dimension API. In this case it would be inserting the schools in the school dimension table.

Time as a dimension

Example of a time dimension

1{ 2 "name": "Last 7 Days", 3 "type": "time", 4 "storage": { 5 "retention": "30 days", 6 "bucket_size": "7 days", 7 } 8}

Here the retention is the time for which the data will be stored in the database and the bucket_size is the time interval for which the data will be aggregated.

cQube has an upper limit to the amount of data that can stored based on the retention policy - 30M records. This is not enforced at the storage layer but the system is not optimized for datasets larger than this. This is kept to keep the architecture simple and not adding an additonal layer of archival storage complexity. This is a soft limit for now and can be increased in the future.

If you have a requirement to increase this limit, please reach out to us. The lowest bucket_size that cQube would support would be a day.

Impact of Dimensions on Datasets

  1. Impact on Datasets

    • cQube Datasets are optimized for query based on dimensions. Indexes are defined as part of the schema.

  2. Impact on Input Data

    • The input data is validated against the dimension data before being persisted.

    • The event schema should always reference a dimension for this to work.

Notes:

  1. An update to the schema is not supported yet. To update please create a new schema and tie it up to the Dataset Model. Please note that this is a slower process and would lead to a downtime of the cQube instance until the migration is completed. Currently, it can only happen through private APIs exposed to a developer.

    1. If this needs to absolutely happen, create a new dimension, insert the dimension data and run the entire event source through the transformers to update the datasets.

  2. A relational database is used to store the data in a flat manner. All fields are flattened and used as column names when creating a dimension table.

  3. The schema provided to insert data is used to validate the data before inserting it into the database. This is to ensure that the data is in the correct format and that the data types are correct. AJV is used to validate the data.

4. Datasets

Dataset insert requires an input that looks like the following.

1{ 2 "name": "attendance_school", 3 "data": { 4 "date": "2019-01-01", // This is a dimension 5 "school_id": 901, // This is a dimension 6 "grade": 1, 7 }, 8 "counter": 190 9} 10

Note: id, created_at, updated_at, count, sum, and percentage are all automatically added to the dataset. Mapping a dataset to a dimension happens the DatasetGrammar which is defined ahead of time.

Dataset Grammar

Defines

  1. How the dataset is defined

  2. How it is linked to dimensions

  3. Which fields need to be cached/indexed to be queried faster

1{ 2 "indexes": [["date"], ["school_id"]] // For faster queries 3 "title": "attendance_school_grade", //name of the dataset 4 "schema": { 5 6 // Will be a JSONSchema. Right now just sharing an example here 7 "date": "2019-01-01", // Dimension 8 "grade": 1, // Addtional filterable data => indexed 9 "school_id": 901, // Dimension 10 }, 11 "dimensionMappings": [ 12 { 13 "key": "school_id", 14 "dimension": { 15 "name": "Dimension1", 16 "mapped_to": "Dimension1.school_id" //Optional in case of time. 17 } 18 }, 19 { 20 "key": "date", 21 "dimension": { 22 "name": "Last 7 Days" 23 } 24 } 25 ] 26}

5. Instruments

cQube implements Opentelemetry Instruments which can be of the following types. A single instrument allows for a lot of functionalities downstream so please ensure you have all the use cases covered before adding a new instrument

  • counter (def) - allows for aggregations as sum, count, percentage, min, max, avg, percentile

  • Future instruments to be implemented - meter, gauge, histogram.

6. Event Management

  1. All events are stored as CSV files in S3. They can be rerun to repopulate Datasets and Dimensions.

  2. Time dimensions-based Datasets are rolled up and bucketed using an OLAP extension in PSQL. This provides the best of both worlds - faster inserts (for data with an older timestamp compared to OLAPs) and almost comparable OLAP rolling aggregations on time dimensions.

Assumptions

  1. Since it’s not a frequent use case, Dimension updates and deletion are not supported. In case this is needed they can just create a new dimension, update the mapping to datasets and refresh the datasets.

  2. Event Sourcing is used throughout the design.

Next Steps

  1. SDKs can be added to push events to Opentelemetry Receivers

  2. The current collection service can be deprecated to use Opentelemetry Receivers.

  3. Transformers can be moved to an Opentelemetry Transformer.

  4. The collector can send data to any destination using official plugins.