# Welcome to Obsrv

Obsrv was conceived and incubated within the open-source initiative [Sunbird](https://sunbird.org/) in 2016. Obsrv is built to operate with highest levels of reliability with bare minimum operations effort at scale.

One of the early and significant applications of Obsrv is within **DIKSHA**, the national platform for schools and teachers in India. In the DIKSHA environment, Obsrv handled data volumes ranging from a few million events per day to **2 billion events per day** at peak. The current average volume is 500 million events per day and the system has not experienced any downtime in the past 2 years.

This microsite provides a comprehensive overview and guides to understand and use Obsrv for various data use cases.

## Getting Started

<table data-view="cards"><thead><tr><th></th><th></th><th data-hidden data-card-cover data-type="files"></th><th data-hidden data-card-target data-type="content-ref"></th></tr></thead><tbody><tr><td><strong>Introduction</strong></td><td>Intro to Obrsv</td><td><a href="/files/9JJowoaQ4P352rvasNSD">/files/9JJowoaQ4P352rvasNSD</a></td><td><a href="/pages/eTCNDdzHhCcpX5mIwFbV">/pages/eTCNDdzHhCcpX5mIwFbV</a></td></tr><tr><td><strong>Core Concepts</strong></td><td>Understand Obsrv Fundamentals</td><td><a href="/files/AZyZcxUPIvxNFlb5iyuN">/files/AZyZcxUPIvxNFlb5iyuN</a></td><td><a href="/pages/eG0hIG7sAJU7aTuamKGs">/pages/eG0hIG7sAJU7aTuamKGs</a></td></tr><tr><td><strong>Case Studies</strong></td><td>Real World Use cases</td><td><a href="/files/aF8cPdEJXSqjoSZXZRBK">/files/aF8cPdEJXSqjoSZXZRBK</a></td><td><a href="/pages/7pvDv04rdGhz3k46RgT7">/pages/7pvDv04rdGhz3k46RgT7</a></td></tr></tbody></table>

## Guides

<table data-view="cards"><thead><tr><th></th><th data-hidden data-card-target data-type="content-ref"></th><th data-hidden data-card-cover data-type="files"></th></tr></thead><tbody><tr><td><strong>Installation</strong></td><td><a href="/pages/UiZZR36UP9ygFuDMRpDx">/pages/UiZZR36UP9ygFuDMRpDx</a></td><td><a href="/files/ELnsiXajMvYgd6UIxs9q">/files/ELnsiXajMvYgd6UIxs9q</a></td></tr><tr><td><strong>APIs</strong></td><td><a href="/pages/k6LM7p39OABQ0LO6QE2Q">/pages/k6LM7p39OABQ0LO6QE2Q</a></td><td><a href="/files/9R46MNUzTFyP7dtOZ9pt">/files/9R46MNUzTFyP7dtOZ9pt</a></td></tr><tr><td><strong>Developer Guide</strong></td><td><a href="/pages/pvr1EzvS3Toy0wVb19TN">/pages/pvr1EzvS3Toy0wVb19TN</a></td><td><a href="/files/vI2neP6BGyRwF9VU0FGs">/files/vI2neP6BGyRwF9VU0FGs</a></td></tr></tbody></table>


# The Value of Data

Importance of Data and its use-cases

In today's digital age, data has become the cornerstone of innovation, driving decision-making processes, and revolutionizing industries across the globe. Below are a few data-driven approaches that drive significant organizational growth:

* **Informed Decision-Making**: Data empowers organizations to make informed decisions backed by evidence rather than intuition alone. By analyzing patterns, trends, and correlations within data sets, businesses can optimize operations, mitigate risks, and capitalize on emerging opportunities.
* **Enhanced Customer Experiences**: Understanding customer behavior through data analytics enables businesses to tailor products, services, and marketing campaigns to meet evolving consumer preferences. Personalization based on data insights fosters stronger customer relationships and increases brand loyalty.
* **Predictive Capabilities**: Advanced analytics and machine learning algorithms leverage historical data to predict future trends, behaviors, and outcomes. By anticipating market shifts, demand fluctuations, and customer needs, organizations can stay ahead of the curve and adapt proactively to changing circumstances.

The value of data lies not in its abundance but in its transformative potential to drive innovation, inform decision-making, and create tangible value across diverse industries. By embracing data-driven methodologies and harnessing the power of advanced analytics, organizations can unlock new insights, capitalize on emerging opportunities, and chart a course towards sustainable growth and success in the digital era.

> ### *<mark style="color:blue;">What's fundamentally needed is a robust “</mark><mark style="color:blue;">**Data Value Chain**</mark><mark style="color:blue;">” capable of unlocking the full potential of an organization's data assets and driving sustainable value creation.</mark>*


# Data Value Chain

Stages of a Data Value Chain & typical implementation

The Data Value Chain represents the journey of data from its raw form to actionable insights and strategic decision-making. The data value chain typically encompasses the stages shown in the below diagram:

<figure><img src="/files/EFtbGwFjn725sItByKyO" alt=""><figcaption><p>Stages of Data Value Chain</p></figcaption></figure>

Efficient management of the Data Value Chain is crucial for all companies seeking to derive maximum value from their data. The figure below shows the typical technology components used in implementing a data value chain.

<figure><img src="/files/LvRmCeduABgWsUdkNktQ" alt=""><figcaption><p>Typical Technology Components in a Data Value Chain</p></figcaption></figure>

> ### *<mark style="color:blue;">**There is growing demand for a robust data value chain in organizations seeking to extract maximum value from their data…**</mark>*


# Challenges

Challenges in implementing and managing a robust Data Value Chain

Despite the availability of numerous scalable and dependable technologies in the data space, the combination of these technologies often results in a fragile end solution. Only select big-tech companies have successfully mastered the processes of generating, consuming, and utilizing data reliably at any scale—examples include Google, Facebook, Netflix, Amazon, and LinkedIn. The existence of reliable data platforms facilitating a robust Data Value Chain has served as a significant distinguishing factor for these companies. And none of these companies have shared their end-to-end solutions with others, only a handful have released some of the tools they internally use, such as Facebook's contribution with Cassandra.

> ### *<mark style="color:blue;">A pressing issue for many organizations is the substantial effort required to develop, operate or maintain an end-to-end data solution reliably.</mark>*

The challenge of reliably managing the data value chain is growing for numerous companies, particularly as contemporary products generate substantial data volumes, even with a relatively modest user base.

### Key Challenges

The challenges faced by most of the organizations in operating a data value chain mainly falls into one of the following four categories:

1. **Time**: Significant amount of time mis-spent in managing the solution rather than in leveraging data's full potential
2. **Cost**: High upfront CapEx and running costs
3. **Capability**: Challenges to build, manage and operate complex data technologies & systems
4. **Risk**: Business & technical risks due to propietary & fragmented solutions, rigid & less reliable systems

Listed below are some challenges that are often faced by organizations with data platforms & solutions:

* Comprehensive, ready-to-use solutions for implementing the entire data value-chain are scarce. In many instances, organizations resort to employing extensive teams of Data and DevOps engineers to build these solutions.
* Data Analysts consistently grapple with challenges related to data integrity, quality, and accessibility, primarily due to the dynamic and evolving nature of data.
* Data Engineers spend substantial time resolving reliability issues due to the agile nature of data and interoperability challenges between components of data platforms.
* Inherent complexities of data platforms, when exposed, increase timelines for new pipeline creation, increasing the lead times for generating data insights.
* Organizations find themselves compelled to transition to a new data solution as they grow and have diverse data use-cases.
* Majority of existing managed solutions, if not all, bind users to proprietary data tools and formats, creating a vendor lock-in.


# The Solution: Obsrv

Obsrv: A Resilient and Reliable Data Value Chain Orchestrator

The ideal solution to address the challenges in creating & operating a data value chain should have the following characterstics:

* Retrieve data from diverse sources, comprehend various formats, and adjust to any changes.
* Handle and store data without requirement of scripting or coding.
* Facilitate data utilization for all scenarios.
* Function reliably at any scale without necessitating modifications.
* Orchestrate the optimal data value chain...

Obsrv brings together the best data tools and technologies, and seamlessly orchestrating their integration through extreme automation techniques. The outcome is an end-to-end low-code data platform that is not only reliable but also resilient across diverse data requirements. Over the course of seven years, Obsrv has evolved to effectively tackle a broad spectrum of challenges in data analysis and data engineering.

> ### *<mark style="color:blue;">Functioning as an orchestrator of the Data Value Chain, Obsrv connects data to its inherent value.</mark>*

Key Benefits of Obsrv include:

* **Unified Data Infrastructure:** Obsrv is an end-to-end low-code data platform, facilitating integration of data from diverse sources, adaptable to changes and generating valuable insights.
* **Built-in Observability:** Obsrv possesses innate observability, knows when the data breaks and avoids data down times, ensuring a continuous and reliable data flow.
* **Reliability by Design:** Obsrv is engineered with extreme automation, ensuring seamless & reliable operation irrespective of the scale at which it is deployed.
* **Instant Data Utilization:** Obsrv facilitates configuration for input data sources, transformations, and the entire pipeline effortlessly, without the need for coding.
* **Diverse Applications:** Obsrv empowers use of data across a spectrum of scenarios, including real-time applications, ensuring organizations stay ahead in the era of rapid data-driven decision-making.
* **Freedom:** Obsrv uses open technologies & formats and its core engine is fully open source, guaranteeing zero lock-in and complete freedom to operate & exit.


# Core Concepts

A deep-dive into core concepts of Obsrv

This section covers key concepts, technical architecture and monitoring capabilities of Obsrv.


# Obsrv Overview

How & Where does Obsrv fit in the data landscape

As touched upon in the introduction section, A data value chain consists of ingesting & processing of data via data pipelines, storage of the processed data in a data-warehouse or data lake, and querying of the data for analytical purposes. The querying of data is either via a batch request or real-time depending upon the underlying storage layer configured.

Therefore many tools and technologies have come up in this space (see diagram below) trying to solve very specific problems as listed below:

1. **Data Integration Platforms:** Data integration platforms (and tools) are used to move the data from operational sources (like OLTP databases, object stores, log streams) be it structured, semi-structured or unstructured into a data store where further processing and querying can happen. Some of the integration platforms also provide the ability to transform the data while moving into the data platform. This movement is either batch or streaming depending on the sources themselves.
2. **Data Warehouses:** Data warehouses are used to store the data in the format that is friendly for analytical queries. Typical analytical queries crunch large amounts of data on any dimension depending on user needs. Typical OLTP databases cannot support these kinds of adhoc and interactive querying needs.
3. **Lake Houses:** While data warehouses are present for storage, most of them support only structured data. In addition, data warehouses have strong schema affinity which make them very slow to adapt for changing needs. For ex: what if new attributes are added to the data? It is huge engineering work to prepare them to be stored in a data warehouse. In addition data warehouses are not AI/ML friendly or efficient. Lake-houses are an evolved architecture pattern to handle the limitations of data warehouses while providing cheap storage (as they are built on top of data lakes) and efficient and fast querying capabilities to AI/ML algorithms.
4. **Full-stack solutions:** While there are many tools, to realize an end-to-end data value chain, many tools have to be stitched together to get an usable data platform. While the tools are scalable by themselves, architecturally ensuring reliability of each tool is a challenge and when stitched together the complexity increases exponentially. There are many full-stack solutions that have tried to solve this problem by providing all the capabilities of data integration tools, data warehouses and lake-houses.

While full-stack solutions themselves offer a complete solution, almost all of them are not built as real-time solutions ground up and are neither open-sourced.

As explained in the diagram below, this is the reason why Obsrv has been built - fully open source, stitching together the best tools for pipelines, storage and querying and real-time first by design.

<figure><img src="/files/M1errZB6jBrtpRoJVMde" alt=""><figcaption><p>Capability Landscape</p></figcaption></figure>


# High Level Architecture

Under the hood of Obsrv

Obsrv fuses multiple technologies together with extreme automation and detailed monitoring coupled with intelligent services to work on any cloud to enable multiple data use-cases via decoupled integrations. The chaining together of these layers give Obsrv its scalability, reliability and efficiency.

<figure><img src="/files/eiWuMje4XPHBIKBD0F2Z" alt=""><figcaption><p>Obsrv in a box - Fuse together multiple layers</p></figcaption></figure>

Following diagram explains the high level architecture of the Obsrv data platform

<figure><img src="/files/ubWLJyN0rVOXUWOyKbef" alt=""><figcaption><p>Obsrv Data Platform</p></figcaption></figure>

Following are the key components in Obsrv:

1. **Connectors:** A connector (that can be literally dropped in) has the ability to pull the data from any source either as a stream/event or batch. The connector framework of Obsrv allows one to develop a connector quickly within a couple of days using popular technologies like Apache Spark and Apache Flink and in the language of their choice - java/scala/python. By design the framework takes care of scaling and reliability of the connectors
2. **Data Pipeline:** Data pipelines are Apache Flink and Apache Spark based jobs that are designed to process data at real-time speed. The data pipeline of Obsrv is extremely elastic and scales from 1 cpu to many cpus with minimal configuration changes and is also customizable and/or extendable.
3. **Real-time OLAP Store:** If configured, all the data is persisted in a real-time OLAP store called Apache Druid that would drive all real-time use-cases.
4. **Data Lake and LakeHouse:** As a default configuration, all data is persisted into a data lake (like S3 object store) and a LakeHouse called Apache Hudi. The LakeHouse and data lake drive the exhaust, AI/ML queries, batch aggregate queries and reporting needs.
5. **Unified Query Engine:** The unified query engine component takes care of all data driven use-cases. It allows the user to query using JSON, SQL and Spark/Trino interfaces.


# Key Capabilities

Key Features to orchestrate the data value chain

### Functional Capabilities

Following are the key functional capabilities of Obsrv across the data value chain. All the capabilities highlighted with <mark style="color:blue;">blue</mark> are available and <mark style="color:orange;">orange</mark> will be available in future releases

| Ingestion                                                                                                                                                                                                                                                                                                                                                                                 | Processing                                                                                                                                                                                                                                                                                                                                                                                                                                | Storage                                                                                                                                                                                                                                                                                                                                                                                                                         | Querying                                                                                                                                                                                                                                                                                                                                                                                                                             |
| ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ |
| <p><mark style="color:blue;">Connectors Framework</mark></p><p><mark style="color:blue;">Connectors Library</mark></p><p><mark style="color:blue;">CSV, JSON Data Formats</mark></p><p><mark style="color:orange;">Parquet, Avro, Protobuf Dataformat</mark></p><p><mark style="color:orange;">Auto-Schema Detection</mark></p><p><mark style="color:orange;">Schema Evolution</mark></p> | <p><mark style="color:blue;">Schema Agnostic</mark></p><p><mark style="color:blue;">Multiple Validation Modes</mark></p><p><mark style="color:blue;">Deduplication</mark></p><p><mark style="color:blue;">Denormalization</mark></p><p><mark style="color:blue;">Custom Transformations</mark></p><p><mark style="color:orange;">Masking & Encryption</mark></p><p><mark style="color:orange;">Jsonata and SQL Transformations</mark></p> | <p><mark style="color:blue;">Multiple Storage Types</mark></p><p><mark style="color:blue;">Data Lake &</mark> <mark style="color:orange;">LakeHouse</mark></p><p><mark style="color:blue;">Real-time OLAP Storage: Optional</mark></p><p><mark style="color:blue;">Archival & Retention Policies</mark></p><p><mark style="color:blue;">Data Exhausts</mark></p><p><mark style="color:orange;">Right to be forgotten</mark></p> | <p><mark style="color:blue;">Aggregate tables</mark></p><p><mark style="color:blue;">SQL, JSON, & Spark Interfaces</mark></p><p><mark style="color:blue;">Time Zone Configuration</mark></p><p><mark style="color:blue;">Geo-spatial Queries</mark></p><p><mark style="color:orange;">Data Aliases</mark></p><p><mark style="color:orange;">Query Access Control</mark></p><p><mark style="color:orange;">Sink Connectors</mark></p> |

### Infra Capabilities

And following are the capabilities provided by the overall infra:

| Infra Capabilities                                                                                                                                                                                                                                                                                                                                                                                                                                                                          |
| ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| <p><mark style="color:blue;">Complete Monitoring</mark> <mark style="color:orange;">Alerts & Notifications</mark> <mark style="color:blue;">Infra Operations</mark></p><p><mark style="color:blue;">Dataset Management</mark> <mark style="color:orange;">Connectors Management</mark> <mark style="color:blue;">System Configuration</mark></p><p><mark style="color:blue;">One-click Install Backup and Restore</mark> <mark style="color:orange;">Auto Scaling API Management</mark></p> |


# Datasets

Definition and details about datasets - a key construct in Obsrv

### Introduction

Datasets along with connectors form the core constructs of Obsrv. Anything and everything we do in Obsrv is related to managing a dataset. A dataset can be created with one ore more sources and a dataset can result in creation of one or more tables.

### What does Dataset do?

Most of the time as Engineers, we tend to focus more on technology and current use-cases that we have and the design would be bounded. Hence with data platforms, we try to solve by thinking about the technologies and designs that can enable ingestion, processing, querying and storage and the use-cases at hand that we have. But as a open-source and future proof product - How do we design for evolving needs?

With Obsrv - we took the route of data first and data centric design and at the core of it is a **Dataset**. All the components of obsrv are designed to enable creation, management and consumption of dataset.

With data first thinking - we are able to solve quite a few design problems and helps with explosion of use-cases later:

1. **Data Isolation:** Instead of creating one mega table, datasets enable structuring of data into many datasets further enable data isolation. A dataset can be created per tenant or region. This would help with Data privacy concerns like GDPR where the data has to be stored locally.
2. **Data Governance:** Since a dataset is atomic and can exist indepently, it enables us to govern it more efficiently. Every dataset can have it's own SLA's, Quality metrics, Access Control and derived tables. For ex: Creating a derived table for every tenant can enable one to implement a capability like "Right to be Forgotten" or "Download my data".
3. **Efficient Scaling:** Not all datasets will be of the same size. With the dataset construct - one can theoritically scale a specific dataset, instead of scaling entire Obsrv. Traditional design would scale (or size) for the largest source or peak volume.
4. **Decoupled & Modular:** Each dataset is decoupled with another dataset. One dataset failure doesn't effect another dataset. This enables interesting use-cases like create a dataset on demand from the same source and direct it for adhoc querying without impact production workloads and shutdown the dataset once adhoc analysis is over
5. **Ease of Operations:** At the end of the day, a solution is built once but has to be operated daily. With the dataset construct, even the operations is bifurcated by dataset. This would enable one to focus on critical datasets than all datasets which would result in operational efficiency and reduced operational teams

### Dataset Management

Following lists the high level capabilities available as part of Dataset via dataset management APIs

1. **Create Dataset:** Create dataset by providing or linking sources
2. **Publish Dataset:** Review and Publish a dataset
3. **Manage Dataset:** Perform actions like edit, retire and archive a dataset. Edit of a dataset creates a new version and the changes are propagated to existing dataset on publish
4. **Monitor Dataset:** Monitor dataset operations. Every dataset exposes various metrics like data load, processing speed, queries per second and query throughput
5. **Manage Tables:** Create aggregate tables. Aggregate tables can be created on all or subset of data
6. **Manage Connectors:** Manage source and destination connectors for the dataset


# Connectors

Another key construct of Obsrv which allows it to connect to wide variety of data sources & formats

### Introduction

Connectors are one of the core construct of Obsrv in addition to datasets. While datasets provide a construct to manage your data processing , storage and querying layers, connectors provide a construct to manage your data ingress and outgress.

### What does connectors do?

While Obsrv automatically enables data "push" into the platform there are many use-cases where one has to pull the data from sources. Trying to write custom scripts/jobs to read from many sources and pushing into Obsrv is a complex problem and is prone to quality and reliability issues. "**Connectors**" as a concept is used to solve the problem of pulling data efficiently and reliably. **Connectors** solve quite a few design problems of new age data platforms:

1. **Decoupled:** Enables Obsrv to be decoupled with the data ingestion from source systems. One source system cannot effect the data flowing in from another source system
2. **Source Data Management:** Data quality, volume and lineage can be segregated by source and can be managed independently.
3. **Efficient Scalability:** Only the specific connector that processes large volume needs to be scaled independently rather than the entire data platform
4. **Pluggability:** Just swap out the out of the box connector with your own custom connector or a market-place connector without any impact to the data platform
5. **Extensability:** Future proof where extensibility is guaranteed by design. Your data has a unique source - build your own connector using the connector framework. Connectors not only pull data from sources but can also sync data to choice of your destination (reverse ETL)

Connectors fall into 4 broad categories:

1. **Database:** Any connector pulling data from a OLTP or NoSQL database.
2. **Stream/Event:** Any connector pulling data from streams or event driven systems (like Kafka, RabbitMQ etc). The connectors of this type can process data in real-time.
3. **File:** Any connector pulling data from file systems or object stores like S3, Azure Blob, GCS, MinIO etc
4. **Application:** Custom application specific connector. For ex: A SAP connector to pull data from SAP system.

### Available Connectors

Following are the connectors available out of the box. Connector with <mark style="color:blue;">blue</mark> are available and <mark style="color:orange;">orange</mark> are under incubation and will be available in future releases

| Connector Type | Connector Source                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                     |
| -------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| Database       | <p><mark style="color:blue;">Postgres</mark>, <mark style="color:blue;">MySQL</mark>, <mark style="color:blue;">DB2</mark>, <mark style="color:blue;">MariaDB</mark></p><p><mark style="color:orange;">Oracle</mark>, <mark style="color:orange;">SQL Server</mark>, <mark style="color:blue;">Amazon RDS</mark>, <mark style="color:blue;">Azure SQL</mark>, <mark style="color:blue;">Google Cloud SQL</mark></p><p><mark style="color:orange;">MongoDB</mark>, <mark style="color:orange;">Cassandra</mark>, <mark style="color:orange;">ElasticSearch</mark></p> |
| Stream         | <mark style="color:blue;">Kafka</mark>, <mark style="color:orange;">Postgresql Debezium</mark>, <mark style="color:orange;">MySQL Debezium</mark>, <mark style="color:blue;">Neo4j Transactions</mark>, <mark style="color:orange;">DB2 Debezium</mark>, <mark style="color:orange;">Oracle Debezium</mark>, <mark style="color:orange;">SQL Server Debezium</mark>, <mark style="color:orange;">MongoDB Debezium</mark>, <mark style="color:orange;">Cassandra Debezium</mark>                                                                                      |
| File           | <mark style="color:blue;">AWS S3</mark>, <mark style="color:orange;">Azure Blob Storage</mark>, <mark style="color:orange;">MinIO</mark>, <mark style="color:orange;">Google Cloud Storage</mark>                                                                                                                                                                                                                                                                                                                                                                    |


# Tech Stack

Technologies Powering Obsrv

Obsrv is offering a comprehensive solution that harnesses the capabilities of these technologies while abstracting their complexities, empowering organizations to tap into the combined power effortlessly, reducing the implementation effort significantly.

<figure><img src="/files/iujXDK4TdzGMq1QefGqr" alt=""><figcaption></figcaption></figure>


# Monitoring

One of innate capabilities of Obsrv to empower seamless operations

Obsrv offers a robust monitoring system to monitor the performance and health of the entire system. The infrastructure, application and service components are all configured to emit metrics and logs. Obsrv uses Prometheus as the default metric aggregation system and Grafana Loki as the default log aggregation system from various components.

<figure><img src="/files/LoXpFm9d3cID35aEVk5c" alt=""><figcaption><p>Obsrv Monitoring Stack</p></figcaption></figure>

### Infrastructure Health and Usage <a href="#kjfbsjshwi01" id="kjfbsjshwi01"></a>

Obsrv provides a way to monitor the health and usage of three main aspects of the entire infrastructure, viz. This provides the capability to analyze and intelligently optimize the infrastructure usage according to the workloads in the cluster.

1. CPU usage
2. Disk usage
3. Memory usage

### Observability Segregation <a href="#id-539d3zbuepvu" id="id-539d3zbuepvu"></a>

Obsrv provides an ability to categorize the various infrastructure and application services based on functionality such as Ingestion, Processing, Storage and Querying. This categorization helps in aggregation of metrics at various functional component levels and helps in providing an easy and insightful way for the operations team to understand incidents or performance degradation.

### Full Text Search on Logs <a href="#id-4dku8r653uao" id="id-4dku8r653uao"></a>

Debuggability and traceability are important tools to ensure a non-intrusive way to understand application failures. Obsrv provides an ability to seamlessly search for errors encountered in various applications and services. Obsrv aggregates logs from multiple services using Promtail and Grafana Loki and provides a unified interface on Grafana to perform analysis on logs and understand failures.


# Explore

Additional details about Obsrv


# Roadmap

***

**Description**: A comprehensive overview of the features released, in-progress, and upcoming in Obsrv.

### **Completed Features (1.2.0-RC)**

| **Feature**                      | **Description**                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                             | **Status**                 | **Release**  |
| -------------------------------- | --------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -------------------------- | ------------ |
| 📦 **Connector Framework**       | <p>- Add and manage connectors via APIs<br>- Create and drop streams (Flik) or batch jobs (Spark) via APIs<br>- Add and manage custom jobs via APIs<br>- <strong>Enhanced stability</strong> with rigorous testing and fixes for key issues<br>- Open-source connectors for <strong>Kafka</strong>, <strong>JDBC</strong>, and <strong>Object Store</strong></p>                                                                                                                                                                                                                                                            | ✅ **Completed**            | **1.2.0-RC** |
| 🔌 **Open-Source Connectors**    | <p>- Supported connectors:<br>• <strong>Kafka</strong><br>• <strong>JDBC</strong><br>• <strong>Object Store</strong></p>                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                    | ✅ **Completed**            | **1.2.0-RC** |
| 🏞️ **Lakehouse**                | - Transactional Data Lake Platform for scalable and reliable data processing and storage                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                    | ✅ **Completed**            | **1.2.0-RC** |
| ⚙️ **Automation Refactor**       | <p>- Stabilized scripts for <strong>greater reliability</strong> and ease of use<br>- Support for <strong>AWS</strong>, <strong>GCP</strong>, <strong>Azure</strong>, and on-premises environments<br>- Streamlined installation for seamless deployments across any setup</p>                                                                                                                                                                                                                                                                                                                                              | ✅ **Completed**            | **1.2.0-RC** |
| 📊 **Dataset Management**        | <p>- Comprehensive dataset management with advanced features:<br>• Automatic schema evolution at processing/storage layers<br>• Seamless import/export<br>• Auto schema generation and indexing<br>• Proactive alerts and monitoring<br>• Masking and encryption<br>• JSONata and SQL transformations<br>• Data exhausts<br>• Data aliases for simplified replay and migration without downtime</p>                                                                                                                                                                                                                         | ✅ **Completed**            | **1.2.0-RC** |
| 🛡️ **RBAC**                     | - Manage users and roles via APIs for robust role-based access control                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                      | ✅ **Completed**            | **1.2.0-RC** |
| 📡 **Obsrv Exporter**            | - Export OpenTelemetry-compliant monitoring data for seamless integration with external monitoring systems                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                  | 🔄 **Partially Completed** | **1.2.0-RC** |
| 🖥️ **Obsrv Management Console** | <p>- Full-featured <strong>Obsrv Management Console</strong> for managing datasets and connectors<br>- Effortless <strong>dataset lifecycle management</strong> (Draft → Ready for Publish → Published → Retired)<br>- <strong>Source Connector Configuration</strong> for <strong>Kafka</strong>, <strong>JDBC</strong>, <strong>Object Store</strong><br>- Seamless <strong>Export/Import</strong> functionality<br>- Automatic <strong>Schema Generation & Indexing</strong><br>- Proactive <strong>Alerts & Monitoring</strong> for dataset health<br>- <strong>Resource Monitoring</strong> for system performance</p> | ✅ **Completed**            | **1.2.0-RC** |
| 📝 **Query APIs**                | <p>- <strong>Query from Data Sources</strong>: Query from multiple data sources such as <strong>Druid</strong>, <strong>Hudi</strong>, <strong>Lakehouse</strong>, and others.<br>- <strong>Exhaust API</strong>: API for retrieving data exhausts from blob stores for reporting and analysis.<br>- <strong>Template API</strong>: API for creating and using query templates, enabling reusable and efficient queries across different datasets and environments.</p>                                                                                                                                                     | ✅ **Completed**            | **1.2.0-RC** |

***

### **Upcoming Features**

| **Feature**                              | **Description**                                                                                                         | **Status**    | **Release** |
| ---------------------------------------- | ----------------------------------------------------------------------------------------------------------------------- | ------------- | ----------- |
| ⚙️ **Stabilization Enhancements**        | - Improve automation scripts for **greater reliability** and ease of use.                                               | ⏳ **Planned** | **Future**  |
| 🔒 **Security & Performance Upgrades**   | - Software updates for critical components like **Druid**, **Flink**, and **Kubernetes**.                               | ⏳ **Planned** | **Future**  |
| 🖥️ **User Experience Improvements**     | - Enhanced management console for a better user experience, simplifying workflows for dataset and connector management. | ⏳ **Planned** | **Future**  |
| 🔗 **Connector Framework Stabilization** | - Optimization and stabilization of the connector framework and its supported connectors.                               | ⏳ **Planned** | **Future**  |

***


# Case Studies

Real-world scenarios demonstrating the applications of Obsrv


# Agri Climate Advisory

Climate variability is a major source of risk to food production as well as to the livelihoods of small and marginal farmers. With other biophysical, socio-economic, and political factors, climate risk contributes enormously to food insecurity, economic losses, and poverty. Climate information services (CIS) can help farmers and food systems mitigate some of these risks as well as build resilience to climate shocks. Research indicates that farmers place value on location and crop specific weather and climate-based agriculture advisories that support farmers in making key decisions while minimizing climate and weather induced risks. Agri-meteorologists play a pivotal role in interpreting climate data and translating it into contextual actionable insights for farmers. However, the effectiveness of their advisories heavily relies on the quality, accessibility, and timeliness of data. Receiving climate information, especially of impending weather events, can help agri-meteorologists recommend pre-emptive actions to farmers and minimize their crop production losses from weather and climate induced events.

In an era where climate change poses unprecedented challenges to agriculture, harnessing data-driven insights becomes crucial for sustainable and productive farming practices. This case study explores how OBSRV enables agri-meteorologists to deliver context-rich advisories to farmers, enhancing agricultural resilience and productivity.

### **The Challenge** <a href="#gerzblv4zvzz" id="gerzblv4zvzz"></a>

Traditionally, accessing and processing diverse climate and soil data from multiple sources posed significant challenges:

* **Data Heterogeneity**: Climate and soil data often originate from disparate sources, such as meteorological departments, global weather datasets, satellite imagery, and soil databases, presenting in different formats.
* **Spatial and Temporal Discrepancies**: Data collected from various sources may exhibit spatial and temporal discrepancies, complicating the harmonization process.
* **Accessibility and Interpretation**: Agri-meteorologists require easily accessible and interpretable data to formulate relevant advisories for farmers.

### **Use Case: Climate-Based Agriculture Advisory** <a href="#maaelwef91q6" id="maaelwef91q6"></a>

In the context of climate-based agriculture advisory, OBSRV facilitates the following processes:

* **Data Acquisition**: OBSRV ingests climate and soil data from multiple sources such as meteorological departments, global weather datasets, satellite imagery providers, and soil databases.
* **Data Harmonization**: The platform harmonizes diverse datasets, aligning them spatially and temporally to ensure consistency.
* **Data Analysis**: Agri-meteorologists leverage OBSRV's analytical capabilities to analyze climate trends, identify risk factors, and forecast agricultural outcomes.

Based on data insights, agri-meteorologists craft context-rich advisories and disseminate them to farmers through various channels, including mobile applications, SMS alerts, and web portals.

### **Benefits** <a href="#ovc3tzj5wsr" id="ovc3tzj5wsr"></a>

The adoption of OBSRV for climate-based agriculture advisory yields the following benefits:

* **Enhanced Decision-Making**: Agri-meteorologists make informed decisions backed by data-driven insights, optimizing agricultural practices and mitigating risks.
* **Improved Resilience**: Farmers receive timely and relevant advisories, enabling them to adapt to changing climatic conditions and minimize yield losses.
* **Resource Optimization**: Precise advisories help farmers optimize resource utilization, including water, fertilizers, and pesticides, leading to cost savings and environmental sustainability.
* **Knowledge Exchange**: OBSRV fosters knowledge exchange and collaboration among agri-meteorologists, researchers, and farmers, promoting innovation and resilience in agriculture.
* **Scalability and Flexibility**: OBSRV is designed to scale according to evolving data needs and can accommodate new & additional data sources and functionalities.

OBSRV is a transformative solution for climate-based agriculture advisory, empowering agri-meteorologists to harness the power of data for informed decision-making. By facilitating the ingestion, harmonization, and accessibility of diverse climate and soil data, OBSRV enables agri-meteorologists to provide context-rich advisories that enhance agricultural resilience, productivity, and sustainability.


# Learning Analytics at Population Scale

In the rapidly evolving landscape of education, analytics has emerged as a pivotal tool for learning platforms at scale, ushering in a transformative era of data-driven insights. As the demand for online learning experiences continues to grow, the importance of analytics in educational platforms becomes increasingly evident. Analytics not only provides a comprehensive view of learner engagement and performance but also empowers educators and administrators with the ability to make informed decisions. Through the systematic analysis of vast datasets, learning platforms can personalize content, identify patterns of student behavior, and optimize the overall learning experience.

This case study explores how OBSRV enabled a population scale learning platform like DIKSHA to provide a means for students to access educational content remotely and offered teachers a digital repository of learning resources and teaching material to facilitate remote instruction, ensuring continuity in learning during the recent pandemic.

## The Challenge

Any learning analytics system which aims at gathering insights into interactions, engagement, and performance to enhance educational outcomes are faced with multifaceted challenges.

* **Scalability**: Ensuring that systems can handle increasing data volumes and user interactions without compromising performance is a persistent challenge.
* **Data Privacy and Security**: Protecting the anonymity of the user data and ensuring proper access controls are present in the system for multi-tenant access.
* **Data Accessibility**: Educators need to have an easy and efficient way of accessing historical data and understand the lineage to come up with specific learning objectives.

<figure><img src="/files/oKUCa9dRxBTqogqY4qKg" alt=""><figcaption><p>Learning Analytics at Scale</p></figcaption></figure>

## Benefits

* **Reliability at Scale**: Educators can make informed decisions in real-time and change the learning plan as Obsrv ensures reliability at scale while processing huge volumes of data.
* **Lossless data processing**: Obsrv ensures that the client systems always get an acknowledgment to ensure the data ingestion and processing are lossless.
* **Data Transparency**: Obsrv enabled the teachers to have transparent access to the learning data to understand and personalize learning paths for various students.
* **Scalability**: DIKSHA’s mission was to reach 200 million children with basic learning experiences and Obsrv delivered effortlessly with 5 million Daily Active Users and processing 2 billion data points a day at peak.

\ <br>


# IOT Observations Infra

In agriculture, the prevalent utilization of IoT devices extends to the measurement of weather conditions and soil properties, facilitating informed decision-making for optimizing crop yields. Integration of monitoring and control systems into agricultural machinery, including tractors and planters, enhances the capabilities for refined agricultural analyses. The data generated by both IoT devices and agricultural machinery are synchronized with their designated integration partner organizations. These devices are expected to emit huge volumes of data and it is essential that a data infrastructure capable of efficiently processing large amounts of data and extracting valuable insights is necessary.

## The Challenge

* **Interoperability**: Integration challenges arise due to the diverse range of IoT devices and sensors from different manufacturers. Ensuring seamless communication and compatibility between devices can be a significant hurdle.
* **Data Management**: Handling large volumes of data generated by IoT devices necessitates effective data management strategies. Storage, processing, and analysis of this data must be efficient to derive meaningful insights without overwhelming the system.
* **Scalability**: As the number of IoT devices increases, managing the scalability of the infrastructure becomes challenging. Scaling up to accommodate a growing network of sensors and devices without compromising performance is a constant concern.

## Use Case: Precision Farming

Obsrv helps the process of precision farming through

* **Data Interoperability**: Obsrv helps in configuration driven transformations to standardize the data across multiple IoT manufacturers/organizations and agricultural machinery.
* **Data Accessibility**: Obsrv provides a simpler way to access/query the data in real-time so that the farmers or field planters/operators can take decisions in real-time.
* **Spatial Aggregations**: Obsrv provides a seamless way to perform spatial aggregations of various measurements across large areas of agricultural fields. For example, a field operator can observe in real-time the amount of pesticide sprayed in a specific portion of the agricultural field assigned to him. This helps in optimizing the usage of the pesticide according to various requirements.

<figure><img src="/files/oegbcmyIdczhm7ut8GQW" alt=""><figcaption><p>IoT Observations Infrastructure</p></figcaption></figure>

## Benefits

* **Water Efficiency**: OBSRV helps farmers optimize the water usage by understanding the right amount of water to each part of the field based on actual needs.
* **Crop Monitoring**: Obsrv helps in continuous monitoring of crops through the IoT observations which in turn aids the farmers to derive insights into environmental conditions, soil health, and plant growth. This enables timely interventions and adjustments to maximize yield.
* **Predictive Analytics**: IoT measurements contribute to predictive analytics, helping farmers anticipate and mitigate potential issues such as diseases, pests, or adverse weather conditions. This proactive approach enhances crop protection and overall productivity.

\ <br>


# Data Driven Features in Learning Platform

\<to be updated>


# Network Observability

\<to be updated>


# Fraud Detection

\<to be updated>


# Performance Benchmarks

Proof of the pudding for scalability of Obsrv

> <mark style="color:orange;">**Note:**</mark> *<mark style="color:blue;">**This is a work in progress page. Following results are from initial benchmarks. Detailed benchmarks will be added shortly once the benchmark exercise is completed**</mark>*

### Cluster Size

| Config Name       | Config Value                                |
| ----------------- | ------------------------------------------- |
| Number of Nodes   | 4                                           |
| Node Size         | 4 core, 16 Gb                               |
| PV size           | 1 TB                                        |
| Installation Mode | Obsrv with monitoring and real-time storage |

### Processing Benchmarks

Processing benchmark is independent on number of datasets created, hence the strategy is to test with volume with all configurations enabled. Disabling any configuration is going to improve throughput

#### Configuration 1

1. Dedup turned on
2. De-normalization configured on 2 master datasets
3. Transformations configured on 2 fields
4. Event size of 1 kb

#### Results

| Flink Configuration                     | Events per Min                                 | Events per hour                                | Events per Day                                 |
| --------------------------------------- | ---------------------------------------------- | ---------------------------------------------- | ---------------------------------------------- |
| 1 CPU, 1GB, 1 task slot, 1 parallelism  | \~ 13k \| 13 Mb                                | \~ 750k \| 780Mb                               | \~ 18Million \| 18Gb                           |
| 2 CPU, 2GB, 2 task slot, 2 parallelism  | \~ 30k \| 30 Mb                                | \~ 1.8Million \| 1.8Gb                         | \~ 40Million \| 40Gb                           |
| 4 CPU, 4Gb, 4 task slots, 4 parallelism | <mark style="color:purple;">In Progress</mark> | <mark style="color:purple;">In Progress</mark> | <mark style="color:purple;">In Progress</mark> |

> Note: Many other scenarios with varying flink configurations are under benchmarking and will be updated post completion

### Secor Backups Benchmark

To ensure there is no data loss across obsrv pipeline all data is backuped to object store using S3. Following are the benchmark results of Secor backups in real-time

#### Configuration 1

1. Total Secor processes - 7
2. Total CPU Allocated - 1.5 cpu
3. Event size of 1 kb

#### Results

<table><thead><tr><th width="192">Events per Min</th><th width="187">Events per hour</th><th width="177">Events per Day</th><th>Events per process</th></tr></thead><tbody><tr><td>~ 1.6 Million | 1.6Mb</td><td>~ 100 Million | 100Gb</td><td>~ 2.4 Billion | 2.4Tb</td><td>~ 300 Million | 300Gb</td></tr></tbody></table>

> Note: In DIKSHA we have observed each secor process with 1cpu was able to upload 200Million events (200 Gb) to Azure blob storage

### Druid Indexing Benchmark

Druid indexing benchmark is dependent on number of datasets created and number of aggregate tables. This benchmark is done with minimal configuration only and can actually linearly scale with the number of CPUs provided

#### Minimum Configuration

| Config Name         | Config Value  |
| ------------------- | ------------- |
| Process Name        | Druid Indexer |
| CPU                 | 0.5           |
| Direct Memory       | 2Gi           |
| Heap                | 9Gi           |
| GlobalIngestionHeap | 8Gi           |
| Workers Count       | 30            |
| Pod Memory          | 11Gi          |

#### Results

<table><thead><tr><th width="154">Num of Tables</th><th width="144">Events per Min</th><th>Events per hour</th><th>Events per Day</th></tr></thead><tbody><tr><td>1</td><td>~ 80k | 80 Mb</td><td>~ 4.8 Million | 4.8 Gb</td><td>~ 110 Million | 110 Gb</td></tr><tr><td>2</td><td>~ 40k | 40 Mb</td><td>~ 2.4 Million | 2.4 Gb</td><td>~ 55 Million | 55 Gb</td></tr><tr><td>3</td><td><mark style="color:purple;">In Progress</mark></td><td><mark style="color:purple;">In Progress</mark></td><td><mark style="color:purple;">In Progress</mark></td></tr><tr><td>4</td><td><mark style="color:purple;">In Progress</mark></td><td><mark style="color:purple;">In Progress</mark></td><td><mark style="color:purple;">In Progress</mark></td></tr><tr><td>5</td><td>~ 35k | 35 Mb</td><td>~ 2.1 Million | 2.1 Gb</td><td>~ 50 Million | 50 Gb</td></tr></tbody></table>

> Note: How does the indexing scale when more cpu resources are provided will be added once the benchmark is complete

### Query Benchmark

Similar to processing, query benchmark is dependent on the volume of data but not on the number of datasets (or tables) created. Query performance will increase linearly with the amount of CPU/Memory assigned to the Druid Historical process

#### Minimum Configuration

| Config Name                | Config Value     |
| -------------------------- | ---------------- |
| Process Name               | Druid Historical |
| CPU                        | 2                |
| Direct Memory              | 4608Mi           |
| Heap                       | 1Gi              |
| Pod Memory                 | 5700Mi           |
| Segment Size               | 4.77Gi           |
| No. of rows per segment    | 5000000          |
| processing.numThreads      | 2                |
| processing.numMergeBuffers | 6                |
| Concurrency                | 100              |

#### RAW Table Results

<table><thead><tr><th width="224">Query</th><th width="138">Query Interval</th><th width="120">Throughput</th><th>Response Times (in ms)</th></tr></thead><tbody><tr><td>Group by on Raw Data</td><td>1 Day</td><td>25 r/s</td><td><p>Avg | Min | Max | 90th</p><p>392 | 80 | 686 | 472</p></td></tr><tr><td>Group by on Raw Data</td><td>7 Days</td><td>4 r/s</td><td><p>Avg | Min | Max | 90th</p><p>4933 | 1277 | 8382 | 5154</p></td></tr><tr><td>Group by on Raw Data</td><td>30 Days</td><td><mark style="color:purple;">In Progress</mark></td><td><mark style="color:purple;">In Progress</mark></td></tr></tbody></table>

#### Aggregate (Rollup) Table Results

<table><thead><tr><th width="224">Query</th><th width="138">Query Interval</th><th width="120">Throughput</th><th>Response Times (in ms)</th></tr></thead><tbody><tr><td>Group by on Aggregate Data</td><td>1 Day</td><td><mark style="color:purple;">In Progress</mark></td><td><mark style="color:purple;">In Progress</mark></td></tr><tr><td>Group by on Aggregate Data</td><td>7 Days</td><td><mark style="color:purple;">In Progress</mark></td><td><mark style="color:purple;">In Progress</mark></td></tr><tr><td>Group by on Aggregate Data</td><td>30 Days</td><td><mark style="color:purple;">In Progress</mark></td><td><mark style="color:purple;">In Progress</mark></td></tr></tbody></table>

> Note: Multiple query types with varying interval and historical configuration combinations are being benchmarked actively and results will be updated once the activity is completed.


# Guides

Documentation to install and operate an Obsrv instance

This section contains various guides to help Obsrv adopters and contributors to work with Obsrv. Best efforts are put in to keep the documentation accurate and up-to-date. However, if any gaps or incorrect information is found, request your help in improving the documentation by reporting issues and/or contributing to the documentation.


# Installation Guide

Instructions to install Obsrv in various environments.

Obsrv has complete automated installation scripts to setup Obsrv in various environments. With the provided automation, Obsrv can be installed in any environment within 1 hour (post completion of pre-requisites setup).

This section provides instructions to setup Obsrv on all the widely used cloud platforms and on-prem data centres.


# AWS

This guide provides detailed, step-by-step instructions for installing and configuring Obsrv on AWS, utilizing Terraform, Terragrunt, and Helm.

***

## Infrastructure Requirements

### 1. System Specifications

* **CPU Requirements**:
  * **Minimum**: 19 CPUs.
  * **Optimal Configuration**: 5 nodes with 4 cores each, totaling 80GB of RAM.

The installation package includes both lakehouse and real-time OLAP storage by default. If the lakehouse component is not required, only the real-time OLAP storage can be installed, reducing requirements to **16 CPUs** and **64GB of RAM**.

In this case, we recommend using **2 nodes with 8 cores each**, totaling **64GB of RAM**, by selecting the **`t2.2xlarge`** AWS instance type.

* **Availability Zones**: All instances should be within the same availability zone to minimize cross-zone data transfer costs. The Obsrv installer will automatically create the EKS (Elastic Kubernetes Service) cluster for you.

### 2. Networking Setup

* **CIDR Block**: Use a `/23` CIDR range (512 IPs) for your environment.
  * Example: A VPC with `10.0.0.0/23` provides IPs from `10.0.0.0` to `10.0.1.255`.
* **Subnets**: Ensure subnets are created in all availability zones within your AWS region.

***

## Prerequisites

Before beginning the installation, make sure the following tools are installed on your Linux-based system:

| **Tool**       | **Version**      | **Installation Command**                                                                                                                                                                                      | **Official Documentation**                                                                       |
| -------------- | ---------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ------------------------------------------------------------------------------------------------ |
| **Terraform**  | 1.5.x or earlier | `curl "https://releases.hashicorp.com/terraform/1.5.2/terraform_1.5.2_linux_amd64.zip" -o terraform.zip && unzip terraform.zip && sudo mv terraform /usr/local/bin/ && rm terraform.zip`                      | [Terraform Install](https://developer.hashicorp.com/terraform/install)                           |
| **Terragrunt** | 0.48 or later    | `curl -OL https://github.com/gruntwork-io/terragrunt/releases/download/v0.49.0/terragrunt_linux_amd64 && sudo mv terragrunt_linux_amd64 /usr/local/bin/terragrunt && sudo chmod +x /usr/local/bin/terragrunt` | [Terragrunt Install](https://terragrunt.gruntwork.io/docs/getting-started/install/)              |
| **Helm**       | 3.10.2 or later  | `curl https://get.helm.sh/helm-v3.10.2-linux-amd64.tar.gz -o helm.tar.gz && tar -zxvf helm.tar.gz && sudo mv linux-amd64/helm /usr/local/bin/`                                                                | [Helm Install](https://helm.sh/docs/intro/install/)                                              |
| **AWS CLI**    | 2.10 or later    | `curl "https://awscli.amazonaws.com/awscli-exe-linux-x86_64.zip" -o "awscliv2.zip" && unzip awscliv2.zip && sudo ./aws/install`                                                                               | [AWS CLI Install](https://docs.aws.amazon.com/cli/latest/userguide/getting-started-install.html) |

***

***

## Installation Steps

### 1. Clone the Obsrv Repository

Start by cloning the Obsrv automation repository and checkout to either the latest release tag or `main`.

```bash
git clone https://github.com/Sunbird-Obsrv/obsrv-automation.git
```

### 2. Configure the Kubernetes Cluster

By executing the following commands which will bring up the kubernetes cluster in the AWS environment of configured region.

1. **Navigate to the Configuration Directory**:

   ```bash
   cd ./obsrv-automation/terraform/aws/vars
   ```
2. **Update Configuration Files**:

   * Open `cluster_overides.tf` and modify the configuration values to match your environment.

   ```bash
   building_block = "obsrv"
   env = "dev"
   region = "us-east-2"
   availability_zones = ["us-east-2a", "us-east-2b", "us-east-2c"]
   timezone = "UTC"
   create_kong_ingress = "true"
   create_vpc = "true"
   create_velero_user = "true"
   eks_node_group_instance_type = ["t2.xlarge"] # Choose depending on your requirements by considering the CPU requirements
   eks_node_group_capacity_type = "ON_DEMAND"
   eks_node_group_scaling_config = { desired_size = 5, max_size = 5, min_size = 1 } # Choose depending on your requirements by considering the CPU requirements
   eks_node_disk_size = 100
   ```
3. **Configure S3 for Cluster State**:

   * Open `obsrv.conf` and update your AWS credentials and bucket names.

   ```bash
   AWS_ACCESS_KEY_ID=<your_access_key_id>
   AWS_SECRET_ACCESS_KEY=<your_secret_access_key>
   AWS_DEFAULT_REGION="us-east-2"
   KUBE_CONFIG_PATH="$HOME/.kube/obsrv-kube-config.yaml"
   AWS_TERRAFORM_BACKEND_BUCKET_NAME="obsrv-tfstate"
   AWS_TERRAFORM_BACKEND_BUCKET_REGION="us-east-2"
   ```

### 3. Run the Installation Script

1. **Navigate to the Infra Setup Directory**:

   ```bash
   cd obsrv-automation/infra-setup
   ```
2. **Make the Script Executable**:

   ```bash
   chmod +x ./obsrv.sh
   ```
3. **Run the Installation**:

   * To start the installation, run the script:

   ```bash
   ./obsrv.sh install --provider aws --config ./obsrv.conf --install_dependencies false
   ```

   * If you want the installer to automatically handle dependencies, set `install_dependencies=true`.

### 4. Verify the Cluster

Once the installation completes, verify that your Kubernetes cluster is up and running:

```bash
kubectl get nodes
```

This should show the nodes in your Kubernetes cluster.

***

## Helm Chart Configuration

### 1. Navigate to the Helm Chart Directory

```bash
cd ./obsrv-automation/helmcharts/
```

### 2. Update AWS Cloud Configuration

> **Note:** `global-cloud-values-aws.yaml` is auto-generated by Terraform (the `aws_cloud_values` module) during the `./obsrv.sh install` step in [Installation Steps](#installation-steps), using the values from `cluster_overrides.tfvars` and resources Terraform creates (S3 bucket names, region, IAM role ARNs, Elastic IP, etc). It gets overwritten every time you run `obsrv.sh install` — do not edit it manually. Shown below just so you know what Terraform fills in:

```yaml
global:
  cloud_storage_provider: "aws"
  cloud_store_provider: "s3"
  cloud_storage_region: "<region>"
  dataset_api_cloud_bucket: "<dataset_bucket_name>"
  config_api_cloud_bucket: "<config_bucket_name>"
  postgresql_backup_cloud_bucket: "<backup_bucket_name>"
  redis_backup_cloud_bucket: "<redis_backup_bucket_name>"
  velero_backup_cloud_bucket: "<velero_backup_bucket_name>"
  cloud_storage_bucket: "<storage_bucket_name>"
  hudi_metadata_bucket: "s3a://<hudi_bucket_name>/hudi"
  cloud_storage_config: |
    '{"identity":"<access-key>","credential":"<secret-key>","region":"<region-name>"}'

  storage_class_name: "gp2"
  checkpoint_bucket: "s3://<checkpoint-bucket-name>"
  s3_access_key: "<aws-access-key>"
  s3_secret_key: "<aws-secret-key>"

kong_annotations:
  service.beta.kubernetes.io/aws-load-balancer-type: nlb
  service.beta.kubernetes.io/aws-load-balancer-scheme: internet-facing
  service.beta.kubernetes.io/aws-load-balancer-eip-allocations: "<elastic-ip>"
  service.beta.kubernetes.io/aws-load-balancer-subnets: "<subnet-id>"

service_accounts:
  enabled: true
  secor: eks.amazonaws.com/role-arn: "<role-arn>"
  dataset_api: eks.amazonaws.com/role-arn: "<role-arn>"
  config_api: eks.amazonaws.com/role-arn: "<role-arn>"
  druid_raw: eks.amazonaws.com/role-arn: "<role-arn>" 
  flink: eks.amazonaws.com/role-arn: "<role-arn>" 
  postgresql_backup: eks.amazonaws.com/role-arn: "<role-arn>" 
  redis_backup: eks.amazonaws.com/role-arn: "<role-arn>" 
  s3_exporter: eks.amazonaws.com/role-arn: "<role-arn>" 
  spark: eks.amazonaws.com/role-arn: "<role-arn>" 

velero-backup:
  credentials:
    useSecret: true
    secretContents:
      cloud: |
        [default]
        aws_access_key_id="<aws-access-key>"
        aws_secret_access_key="<aws-secret-key>"

trino:
  additionalCatalogs:
    lakehouse: |-
      connector.name=hudi
      hive.metastore.uri=thrift://hudi-hms.hms.svc:9083
      hive.s3.aws-access-key=<aws-access-key>
      hive.s3.aws-secret-key=<aws-secret-key>
      hive.s3.ssl.enabled=false
```

### 3. Update Domain Configuration

In `global-values.yaml`, replace `<domain>` with your actual domain or Elastic IP:

```yaml
domain: "<domain>.sslip.io"
```

### 4. Install Obsrv

Make the script executable and set the environment variables and run the installation:

```bash
export cloud_env=aws
export AWS_ACCESS_KEY_ID=<aws-access-key>
export AWS_SECRET_ACCESS_KEY=<aws-secret-key>
export AWS_DEFAULT_REGION=<aws-region>
export KUBE_CONFIG_PATH="$HOME/.kube/obsrv-kube-config.yaml"
export KUBECONFIG="$HOME/.kube/obsrv-kube-config.yaml"
chmod +x ./kitchen/install.sh
./kitchen/install.sh core-setup
./kitchen/install.sh all
```

`core-setup` installs the bootstrap CRDs, prerequisites, core database, and Kafka — `all` (migrations, monitoring, oauth, coreinfra, obsrvapis, obsrvtools, additional) depends on these being in place first, so run them in this order.

***

## Post-Installation Verification

After completing the installation, follow these steps to verify that all components are running correctly:

### 1. Check Kubernetes Components

1. **Verify all pods are running**:

   ```bash
   kubectl get pods -A
   ```

   All pods should be in `Running` state. Common namespaces to check:

   * `flink`: Core Pipeline
   * `monitoring`: Monitoring stack
   * `dataset-api`: Dataset APIs
   * `web-console`: Dataset Management console
2. **Check Services**:

   ```bash
   kubectl get svc -A
   ```

   Verify that essential services have external IPs assigned, particularly the Kong service.

If any component fails these checks, refer to the component-specific logs:

```bash
kubectl logs -f <pod-name> -n <namespace>
```

***

## Troubleshooting

### "Unreadable module directory" / `lstat ../modules: no such file or directory` on `terragrunt init`

Terragrunt copies the source directory into `.terragrunt-cache` before running. If `terraform/aws/terragrunt.hcl` has no explicit `source` block, newer Terragrunt versions (`v1.x`) copy only the `aws` directory, breaking the `../modules` relative paths used by `main.tf`. Older Terragrunt versions (`~0.45`) didn't hit this because they ran in place.

**Fix**: ensure `terraform/aws/terragrunt.hcl` has a `terraform { source = ... }` block pointing at the repo root with the `//terraform/aws` working-dir suffix, so both `../modules` and `../../helmcharts` resolve correctly from the cache copy.

### `Failed to resolve provider packages: locked provider ... does not match configured version constraint`

`terraform/aws/.terraform.lock.hcl` can drift out of sync with the provider version constraints in `main.tf` (e.g. `main.tf` requiring `~> 6.0` while the lock file is still pinned to a `4.x` version). This happens when the constraint is bumped in a commit without regenerating the lock file.

**Fix**:

```bash
cd terraform/aws
tofu init -upgrade -backend=false
```

Commit the regenerated `.terraform.lock.hcl`.

### `S3 bucket does not exist` on `terragrunt init`

The remote state bucket configured in `obsrv.conf` (`AWS_TERRAFORM_BACKEND_BUCKET_NAME`) must exist before `terragrunt init` runs. Either create it manually with `aws s3api create-bucket`, or pass `--backend-bootstrap` to auto-provision it (requires Terragrunt with backend-bootstrap support, e.g. `v1.x`):

```bash
terragrunt init --backend-bootstrap
```

### Certificate shows valid in `cert-manager`, but browser still says "Not Secure"

`kubectl get certificate` reporting `Ready: True` only means cert-manager issued the cert into its Kubernetes Secret — it doesn't guarantee every service actually serving TLS on your domain can read that secret. If a service's Ingress lives in a namespace other than where the cert's Secret was created (commonly `web-console`), the `kubernetes-reflector` controller must be explicitly allowed to copy the secret into that namespace via the `reflector.v1.k8s.emberstack.com/reflection-allowed-namespaces` annotation on the Certificate (set in `helmcharts/services/letsencrypt-ssl/templates/issuer-and-certs.yaml`). A namespace missing from that list (e.g. `dataset-api`) causes the Kong ingress controller to fail fetching the secret for that Ingress, which can break Kong's config sync entirely and make it fall back to its default self-signed certificate for **every** host, not just the affected one.

**Diagnose** — check what cert is actually being served on the wire, independent of Kubernetes state:

```bash
echo | openssl s_client -connect <your-ip>:443 -servername <your-domain> 2>/dev/null | openssl x509 -noout -issuer -subject
```

If the issuer shows `O=Kong, CN=localhost` instead of Let's Encrypt, check the Kong ingress controller logs for `failed to fetch the secret` errors:

```bash
kubectl logs -n kong-ingress -l app.kubernetes.io/name=kong -c ingress-controller --tail=50 | grep -i error
```

**Fix**: add the missing namespace to the `reflection-allowed-namespaces` annotation, then re-run the `migrations` helm bundle and restart Kong:

```bash
cd helmcharts/kitchen
cloud_env=aws ./install.sh migrations
kubectl rollout restart deployment kong -n kong-ingress
```

***

## Upgrade Steps

1. **Pull the Latest Code**:

   ```bash
   cd ./obsrv-automation
   git pull
   cd ./infra-setup
   ```
2. **Update Configurations**: Review and update configuration values as needed.
3. **Run Terraform for Upgrade**:

   ```bash
   ./obsrv.sh install --provider aws --config ./obsrv.conf --install_dependencies false
   ```
4. **Upgrade with Updated Cloud Values**:

   ```bash
   export cloud_env=aws
   export AWS_ACCESS_KEY_ID=<aws-access-key>
   export AWS_SECRET_ACCESS_KEY=<aws-secret-key>
   export AWS_DEFAULT_REGION=<aws-region>
   export KUBE_CONFIG_PATH="$HOME/.kube/obsrv-kube-config.yaml"
   export KUBECONFIG="$HOME/.kube/obsrv-kube-config.yaml"
   chmod +x ./kitchen/install.sh
   ./kitchen/install.sh all
   ```

***

By following these steps, you will ensure a successful installation and configuration of Obsrv on AWS.


# Azure

Instructions to install in Azure infra

### Infrastructure Requirements

1. **Minimal Installation:**
   * You need a system with a minimum of 16 CPUs. We recommend using 2 nodes with 8 cores each, totalling 64GB.
   * All machines must reside in the same availability zone to avoid data transfer charges across zones. Our Obsrv installer will automatically create the AKS cluster for you.
2. **Networking Environment:**
   * Ensure your environment has a CIDR of 16
     * For example, a Virtual Network with a CIDR of `10.0.0.0/23` will have IP addresses ranging from `10.0.0.0` to `10.0.1.255`.
   * Subnets must be created in all availability zones within your region.

### Software Prerequisites

Installation of Obsrv requires the following CLI tools as prerequisites. Please note that the following instructions for installing the prerequisites are provided only for **Linux based operating systems**. Please follow the instructions for the specific tools depending upon your operating system.

#### Terraform

* Terraform CLI version 1.5.x or older. Versions above 1.5.x are not MPL licensed.

  <pre class="language-bash" data-overflow="wrap"><code class="lang-bash">curl "https://releases.hashicorp.com/terraform/1.5.2/terraform_1.5.2_linux_amd64.zip" -o "terraform.zip" &#x26;&#x26; unzip terraform.zip &#x26;&#x26; sudo mv terraform /usr/local/bin/ &#x26;&#x26; rm terraform.zip
  </code></pre>
* Download from here - <https://developer.hashicorp.com/terraform/install>

#### Terragrunt

* Terragrunt CLI version 0.48 or later.

  <pre class="language-bash" data-overflow="wrap"><code class="lang-bash">curl -OL https://github.com/gruntwork-io/terragrunt/releases/download/v0.49.0/terragrunt_linux_amd64 &#x26;&#x26; sudo mv terragrunt_linux_amd64 /usr/local/bin/terragrunt &#x26;&#x26; sudo chmod +x /usr/local/bin/terragrunt
  </code></pre>
* Download from here - <https://terragrunt.gruntwork.io/docs/getting-started/install/>

#### Terrahelp

* Terrahelp version 0.7.5 or later

  <pre class="language-bash" data-overflow="wrap"><code class="lang-bash">curl -OL https://github.com/opencredo/terrahelp/releases/download/v0.7.5/terrahelp_0.7.5_linux_386.tar.gz &#x26;&#x26; tar -xzf terrahelp_0.7.5_linux_386.tar.gz &#x26;&#x26; sudo mv terrahelp /usr/local/bin/terrahelp &#x26;&#x26; sudo chmod +x /usr/local/bin/terrahelp
  </code></pre>
* Download from here - <https://github.com/opencredo/terrahelp?tab=readme-ov-file#installation>

#### Helm

* Helm version 3.10.2

  <pre class="language-bash" data-overflow="wrap"><code class="lang-bash">curl https://get.helm.sh/helm-v3.10.2-linux-amd64.tar.gz -o helm.tar.gz &#x26;&#x26; tar -zxvf helm.tar.gz &#x26;&#x26; sudo mv linux-amd64/helm /usr/local/bin/
  </code></pre>
* Download from here - <https://helm.sh/docs/intro/install/>

#### Azure CLI

* Azure CLI tool version 2.10 or later.

  ```bash
  curl -sL https://aka.ms/InstallAzureCLIDeb | sudo bash
  ```
* Download from here - <https://learn.microsoft.com/en-us/cli/azure/install-azure-cli>
* Post installation, authenticate Azure CLI. Please refer to this [link](https://learn.microsoft.com/en-us/cli/azure/get-started-tutorial-1-prepare-environment?tabs=bash#sign-in-to-azure-using-the-azure-cli) for more details about Signing In to Azure CLI

  ```bash
  az login --allow-no-subscriptions
  ```

### Installation Steps:

1. **Clone the `obsrv-automation` repository**:

   ```bash
   git clone https://github.com/Sanketika-Obsrv/obsrv-automation.git
   ```
2. **Navigate to the setup directory**:

   ```bash
   cd ./obsrv-automation/terraform/azure
   ```
3. **Export Azure Credentials**:

   Create a Resource Group, Storage Account and Storage Container through the Azure Portal. Once completed, export the below values as an environment variable.

   ```bash
   export AZURE_TERRAFORM_BACKEND_RG=<resource_group>
   export AZURE_TERRAFORM_BACKEND_STORAGE_ACCOUNT=<storage_account_name>
   export AZURE_TERRAFORM_BACKEND_CONTAINER=<storage_container_name>
   ```
4. **Create the AKS cluster**

   The following commands will create an AKS Cluster

   ```bash
   terragrunt init
   terragrunt  apply -target module.aks -auto-approve
   terragrunt apply -auto-approve
   ```

   During creation of the cluster, you will be asked for prompts as and when required by the installation. Here is a sample of the inputs you have to provide while the above script executes.

   ```bash
   env: dev
   building_block: obsrv
   location: East US 2
   ```
5. Make a note of **Resource Group** created during the cluster creation. Usually it is a combination of `<building_block>-<env>`. For the above example the resorce group will be `obsrv-dev`. You can look for the logs for the statement like below.

   <pre class="language-bash" data-overflow="wrap"><code class="lang-bash">module.network.azurerm_resource_group.rg: Creation complete after 3s [id=/subscriptions/&#x3C;uuid>/resourceGroups/&#x3C;your-resource-group>]
   </code></pre>

   Export the **Resource Group** name as an environment variable

   ```bash
   export AZ_RESOURCE_GROUP=<your-resource-group>
   ```

### Upgrade Steps:

1. Take latest code from **`obsrv-automation`** repository

   ```bash
   cd ./obsrv-automation
   git pull
   cd ./terraform/azure
   ```
2. Ensure all the configuration configured during the installation is properly updated in all places.
3. Run the terraform to upgrade the cluster to the latest versions.

   ```
   cd ./obsrv-automation
   git pull
   cd ./terraform/azure
   terragrunt apply -auto-approve
   ```


# GCP

Instructions to install in GCP infra

### Infrastructure Requirements

1. **Minimal Installation:**
   * You need a system with a minimum of 16 CPUs. We recommend using 2 nodes with 8 cores each, totalling 32GB.
   * All machines must reside in the same zone to avoid data transfer charges across zones. Our Obsrv installer will automatically create the GKE cluster for you.
2. **Networking Environment:**
   * The networking requirements change w\.r.t the type of cluster we are creating. Follow this [link](https://cloud.google.com/kubernetes-engine/docs/concepts/types-of-clusters#availability) to find out different types of clusters supported by Google Cloud in the region you are installation.
   * By default our setup script creates a Single Zonal Public Cluster and the required VPC’s for it.

### Software Prerequisites

Installation of Obsrv requires the following CLI tools as prerequisites. Please note that the following instructions for installing the prerequisites are provided only for **Linux based operating systems**. Please follow the instructions for the specific tools depending upon your operating system.

#### Terraform

* Terraform CLI version 1.5.x or older. Versions above 1.5.x are not MPL licensed.

  <pre class="language-bash" data-overflow="wrap"><code class="lang-bash">curl "&#x3C;https://releases.hashicorp.com/terraform/1.5.2/terraform_1.5.2_linux_amd64.zip>" -o "terraform.zip" &#x26;&#x26; unzip terraform.zip &#x26;&#x26; sudo mv terraform /usr/local/bin/ &#x26;&#x26; rm terraform.zip
  </code></pre>
* Download from here - <https://developer.hashicorp.com/terraform/install>

#### Terragrunt

* Terragrunt CLI version 0.48 or later.

  <pre class="language-bash" data-overflow="wrap"><code class="lang-bash">curl -OL https://github.com/gruntwork-io/terragrunt/releases/download/v0.49.0/terragrunt_linux_amd64 &#x26;&#x26; sudo mv terragrunt_linux_amd64 /usr/local/bin/terragrunt &#x26;&#x26; sudo chmod +x /usr/local/bin/terragrunt
  </code></pre>
* Download from here - <https://terragrunt.gruntwork.io/docs/getting-started/install/>

#### Terrahelp

* Terrahelp version 0.7.5 or later

  <pre class="language-bash" data-overflow="wrap"><code class="lang-bash">curl -OL https://github.com/opencredo/terrahelp/releases/download/v0.7.5/terrahelp_0.7.5_linux_386.tar.gz &#x26;&#x26; tar -xzf terrahelp_0.7.5_linux_386.tar.gz &#x26;&#x26; sudo mv terrahelp /usr/local/bin/terrahelp &#x26;&#x26; sudo chmod +x /usr/local/bin/terrahelp
  </code></pre>
* Download from here - <https://github.com/opencredo/terrahelp?tab=readme-ov-file#installation>

#### Helm

* Helm version 3.10.2 or later

  <pre class="language-bash" data-overflow="wrap"><code class="lang-bash">curl https://get.helm.sh/helm-v3.10.2-linux-amd64.tar.gz -o helm.tar.gz &#x26;&#x26; tar -zxvf helm.tar.gz &#x26;&#x26; sudo mv linux-amd64/helm /usr/local/bin/
  </code></pre>
* Download from here - <https://helm.sh/docs/intro/install/>

#### Setup Google Cloud SDK

* Install Google Cloud SDK by following the instructions [here](https://cloud.google.com/sdk/docs/install)

  <pre class="language-bash" data-overflow="wrap"><code class="lang-bash">gcloud init
  gcloud auth application-default login
  </code></pre>
* Install additional dependencies to authenticate with GKE. Please see [Installing the gke-gcloud-auth-plugin](https://cloud.google.com/kubernetes-engine/docs/how-to/cluster-access-for-kubectl) for reference.

  ```bash
  gcloud components install gke-gcloud-auth-plugin
  ```

#### Create a GCP Project

* Create a project on google cloud and export it as a variable. Please see [Creating and Managing Projects](https://cloud.google.com/resource-manager/docs/creating-managing-projects) for reference.

  ```bash
  export GOOGLE_PROJECT_ID=<myproject>
  ```
* Enable the Kubernetes Engine API for the created project. Please see [Enabling the Kubernetes Engine API](https://cloud.google.com/kubernetes-engine/docs/how-to/creating-a-zonal-cluster#enable-api) for reference.

### Installation Steps:

1. **Clone the `obsrv-automation` repository**:

   ```bash
   git clone https://github.com/Sanketika-Obsrv/obsrv-automation.git
   ```
2. **Navigate to the setup directory**:

   ```bash
   cd ./obsrv-automation/terraform/gcp
   ```
3. **Update Configuration Files**:

   Update the `vars/cluster_overrides.tfvars` file to match your environment and requirement settings. Example configuration in `cluster_overides.tf`

   ```bash
   project                       = "<myproject>"
   building_block                = "obsrv"
   env                           = "dev"
   region                        = "us-central-1"
   gke_cluster_location          = "us-central-1"
   zone                          = "us-central-1-a"
   timezone                      = "UTC"
   service_type                  = "LoadBalancer"

   # cluster sizing

   gke_node_pool_instance_type   = "c2d-standard-8"
   gke_node_pool_scaling_config = {
     desired_size = 2
     max_size = 2
     min_size = 1
   }

   # Image Tags
   command_service_image_tag    = "1.0.0-GA"
   web_console_image_tag        = "1.0.0-GA"
   dataset_api_image_tag        = "1.0.2-GA"
   flink_image_tag              = "1.0.1-GA"
   secor_image_tag              = "1.0.0-GA"
   superset_image_tag           = "3.0.2"
   ```

   Note: if you are changing the region, also check if the instance type is supported in that region. If not please change the instance type from the supported list which can be found in this [link](https://cloud.google.com/compute/docs/regions-zones#available)
4. **GCS Bucket:**

   Obsrv installation requires a GCS bucket to be created which will be subsequently used for storing the cluster state. The bucket name will then be referenced in the installation script mentioned below. For e.g, if the bucket created is **obsrv-cluster-state**, please export it as a environment variable as show below, which will be referenced in the `terragrunt.hcl` under the environment variable **GOOGLE\_TERRAFORM\_BACKEND\_BUCKET**. The bucket will be created automatically if it doesn’t exist. The region can be updated by passing it in an environment variable **GOOGLE\_TERRAFORM\_BACKEND\_BUCKET\_REGION**

   ```bash
   export GOOGLE_TERRAFORM_BACKEND_BUCKET=obsrv-cluster-tfstate
   export GOOGLE_TERRAFORM_BACKEND_BUCKET_REGION=us-central1
   ```
5. **Run Installation Script**: Execute the following command to start the installation process:

   ```bash
   terragrunt init
   terragrunt plan --var-file=vars/cluster_overrides.tfvars # to verify if everything is in order
   terragrunt apply --var-file=vars/cluster_overrides.tfvars
   ```
6. **Monitor Installation Progress**: The script will begin installing Obsrv within the GKE cluster. Monitor the progress and follow any on-screen prompts or instructions.
7. **Completion**: Once the installation process is complete, verify that Obsrv has been successfully installed and configured within your GKE cluster.

### Upgrade Steps:

1. Take latest code from **`obsrv-automation`** repository

   ```bash
   cd ./obsrv-automation
   git pull
   cd ./terraform/gcp
   ```
2. Ensure all the configuration configured during the installation is properly updated in all places.
3. Run the terraform to upgrade the cluster to the latest versions.

   ```bash
   terragrunt apply --var-file=vars/cluster_overrides.tfvars
   ```


# OCI

Instructions to install in OCI

> <mark style="color:blue;">Obsrv installation on OCI is currently under testing. OCI Installation Guide documentation will be published soon...</mark>


# Data Center

Instructions to install in on-prem data centre

### Infrastructure Requirements

1. **Minimal Installation:**
   * You need a system with a minimum of 16 CPUs and 32GB RAM.
   * If running a multi-node cluster, ensure the network reachability between the nodes.

### Software Prerequisites

Installation of Obsrv requires the following CLI tools as prerequisites. Please note that the following instructions for installing the prerequisites are provided only for **Linux based operating systems**. Please follow the instructions for the specific tools depending upon your operating system.

#### MinIO

* Obsrv Installation requires an Object Store for backups, checkpointing and to store other configurations. We have tested our installations with MinIO, and is the recommended one for quick setup. Follow the instructions from <https://min.io/download> to install.

#### Helm

* Helm version 3.10.2 or later

  <pre class="language-bash" data-overflow="wrap"><code class="lang-bash">curl https://get.helm.sh/helm-v3.10.2-linux-amd64.tar.gz -o helm.tar.gz &#x26;&#x26; tar -zxvf helm.tar.gz &#x26;&#x26; sudo mv linux-amd64/helm /usr/local/bin/
  </code></pre>
* Download from here - <https://helm.sh/docs/intro/install/>

#### Addons

* Ensure LoadBalancer is available. For example on a a local setup add-ons such as `metallb` can be enabled and configured. Following is a sample for minikube

  ```bash
  minikube addons enable metallb
  minikube addons configure metallb
  ```
* Ensure metrics is enabled. Following is a sample for minikube

  ```bash
  minikube addons enable metrics-server
  ```

#### Installation Steps:

1. Clone the `obsrv-automation` repository:

   ```bash
   git clone https://github.com/Sanketika-Obsrv/obsrv-automation.git
   ```
2. Navigate to the setup directory:

   ```bash
   cd ./obsrv-automation/terraform/modules/helm/unified-helm
   ```
3. Update the following values in `obsrv/values.yaml` to reflect your MinIO environment.

   ```yaml
   cloud-storage-provider: &global-cloud-storage-provider "s3"
   cloud-storage-region: &global-cloud-storage-region "<MINIO_REGION>"
   s3_bucket: &global-s3-bucket "<MINIO_BUCKET>"
   s3_access_key: &global-s3-access-key "<MINIO_ACCESS_KEY>"
   s3_secret_key: &global-s3-secret-access-key "<MINIO_SECRET_KEY>"
   region: &global-region "<MINIO_REGION>"
   s3_endpoint_url: &global-s3-endpoint-url "<MINIO_ENDPOINT_URL>"
   s3_path_style_access: &global-s3-path-style-access "true"
   ```
4. Export the KUBECONFIG environment variable with for your cluster. For example the below command is to set to its default path

   ```bash
   export KUBECONFIG=~/.kube/kubeconfig.yaml
   ```
5. First create the CRD’s that are required to install Obsrv

   ```bash
   kubectl create -f ./crds/
   ```
6. Run the below command to install the services. The following command may fail a couple of times due to timeouts while downloading the images. Run the same command for a couple of times incase of any errors for the installation to be successful

   <pre class="language-bash" data-overflow="wrap"><code class="lang-bash">helm upgrade --install obsrv obsrv --namespace obsrv --create-namespace --atomic --debug --timeout 3600s
   </code></pre>

### Upgrade Steps:

1. Take latest code from **`obsrv-automation`** repository

   ```bash
   cd ./obsrv-automation
   git pull
   cd ./terraform/modules/helm/unified-helm
   ```
2. Ensure all the configuration configured during the installation is properly updated in all places.
3. Run the helm upgrade the cluster to the latest versions.

   <pre class="language-bash" data-overflow="wrap"><code class="lang-bash">helm upgrade --install obsrv obsrv --namespace obsrv --create-namespace --atomic --debug --timeout 3600s
   </code></pre>


# API Specification


# Dataset Management APIs

List of APIs to manage datasets

### Dataset APIs

Dataset APIs allow you to manage the datasets, such as creating new datasets, updating and list existing datasets and retrieving specific dataset metadata.

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/v2/datasets/create" method="post" expanded="true" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/v2/datasets/update" method="patch" expanded="true" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/v2/datasets/read/{dataset\_id}" method="get" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/v2/datasets/list" method="post" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/v2/files/generate-url" method="post" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/v2/datasets/status-transition" method="post" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/v2/datasets/dataschema" method="post" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/v2/datasets/status-transition" method="post" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/v2/datasets/export/{dataset\_id}" method="post" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/v2/datasets/import" method="get" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/v2/datasets/copy" method="post" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}


# Connector APIs

List of APIs to manage data input connectors

### Connector Config APIs

Connectors are used to ingest data from external sources. It provides a set of APIs to register, list and read connectors configurations API

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/v2/connector/register" method="post" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/v2/connectors/read/{connector\_id}" method="get" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/v2/connectors/list" method="post" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}


# Data In & Out APIs

List of APIs to write and read data

### Data In (Write) API

The Data In APIs facilitate writing data into datasets, designed to simplify the ingestion of JSON data

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/v2/data/in/{dataset\_id}" method="post" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

### Data Out (Read) APIs

Use the Data Out APIs to fetch data from a dataset with support for SQL and Druid native query formats, suitable for analytics and visualization.

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/v2/data/query/{dataset\_id}" method="post" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

### Template APIs

The Template APIs enable efficient management of templates, including creating new templates, updating existing ones, listing all available templates, and retrieving metadata for specific templates. These APIs are designed to facilitate data retrieval from various sources using predefined template formats. They are particularly recommended for production environments where queries remain relatively static and do not require frequent changes.

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/v2/template/create/{template\_id}" method="post" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/v2/template/list" method="post" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/v2/template/read/{template\_id}" method="get" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/v2/template/update/{template\_id}" method="patch" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/v2/template/query/{template\_id}" method="post" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/v2/template/delete/{template\_id}" method="delete" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}


# Alerts and Notification Channels APIs

List of APIs to manage alerts and notification channels

## Notification Channel APIs

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/alerts/v1/notifications/create" method="post" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/alerts/v1/notifications/publish/{alert\_id}" method="get" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/alerts/v1/notifications/search" method="post" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/alerts/v1/notifications/update/{alert\_id}" method="patch" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/alerts/v1/notifications/delete/{alert\_id}" method="delete" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/alerts/v1/notifications/test" method="post" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

## Alerts APIs

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/alerts/v1/create" method="post" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/alerts/v1/publish/{alert\_id}" method="get" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/alerts/v1/get/{alert\_id}" method="get" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/alerts/v1/search" method="post" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/alerts/v1/update/{alert\_id}" method="patch" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/alerts/v1/delete/{alert\_id}" method="delete" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

## Alerts Silience APIs

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/alerts/v1/silence/create" method="post" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/alerts/v1/silence/get/{alert\_id}" method="get" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/alerts/v1/silence/search" method="get" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}

{% openapi src="<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>" path="/alerts/v1/silence/update/{alert\_id}" method="get" %}
<https://raw.githubusercontent.com/Sanketika-Obsrv/obsrv-api-service/main/api-service/swagger-doc/openapi_v2.yml>
{% endopenapi %}


# Dataset Management Console

The **Management Console** is a sophisticated and intuitive platform designed to give you complete control over your datasets, connectors, and system operations. With a user-friendly interface, it streamlines the processes of dataset creation, configuration, and monitoring, while offering advanced features for data transformation, health monitoring, and alert management. Built-in metrics and real-time dashboards provide a comprehensive overview of your system's performance, ensuring efficient and seamless data workflows.

***

## Key Features

* **Dataset Management**: Create, update, and organize datasets with ease.
* **Connector Configuration**: Integrate connectors effortlessly to enable automatic data synchronization between sources and destinations.
* **Data Transformation**: Apply advanced transformations, including masking, encryption, and JSONAT, to ensure data privacy and compliance.
* **Dataset Joining**: Link datasets with master datasets to enrich data analysis capabilities.
* **Duplicate Record Removal**: Ensure data integrity by identifying and removing duplicate records.
* **System Monitoring**: Access real-time dashboards and metrics for monitoring system performance and health.
* **Custom Alerts**: Configure alerts and notifications to proactively manage dataset and connector issues.

***

## Accessing the Management Console

To use the **Management Console**, navigate to the configured **ingress URL** or use port forwarding to access the application. Log in using the default credentials:

* **Username**: `obsrv_admin`
* **Password**: `enDoPvTAxFSd`

Upon logging in, you can explore the extensive features for managing data workflows, transforming data, and monitoring system health.

***

## Key Workflows in the Management Console

### **Datasets Creation**

The Management Console provides a streamlined workflow for creating datasets.

### **1. Connectors**

Manage and configure connectors for seamless data integration.

1. Register connectors via API and view them in the console.
2. Explore default open-source connectors such as JDBC, Kafka, Object Store, and Knowlg.
3. Add custom connectors to meet specific requirements.
4. During dataset creation, select and configure the appropriate connector by filling in the required details.

<figure><img src="/files/VYoRwWxeQfJNyxg0F07L" alt=""><figcaption><p>Connectors List Page</p></figcaption></figure>

<figure><img src="/files/wKcs9yyK1AHi7u7oK4cs" alt=""><figcaption><p>Connector Detailed Page</p></figcaption></figure>

***

### **2. Provide Basic Information**

1. Enter the name of the dataset.
2. Upload a sample data file or schema file if available.

<figure><img src="/files/4yFmpu05FAlmRN9IUVqy" alt=""><figcaption><p>Dataset Creation</p></figcaption></figure>

***

### **3. Ingestion Configuration**

1. Review the schema configuration and make modifications if needed.
2. Resolve any data type conflicts automatically with recommendations provided by the console.

<figure><img src="/files/Ammb0gb0yWyVxHkZRFU0" alt=""><figcaption><p>Ingestion Configuration</p></figcaption></figure>

***

### **4. Processing Configuration**

1. Choose the schema evaluation type.
2. Apply data transformations like masking and encryption for PII fields.
3. Automatically detect and address data type conflicts.
4. Enable settings to prevent duplicate record ingestion.

<figure><img src="/files/WbR2w6AAha8ZFxrs8Zek" alt=""><figcaption><p>Processing Configuration</p></figcaption></figure>

***

### **5. Storage Configuration**

1. Configure storage type, such as lakehouse or real-time OLAP storage.
2. Set up storage configurations, including partition keys, timestamp keys, and primary keys.

Once all sections are configured, review the dataset settings and publish the dataset to make it live.

<figure><img src="/files/B07SryvqmoAOhdJbXioZ" alt=""><figcaption><p>Storage Configuration</p></figcaption></figure>

***

### **6. Dataset Management**

1. View all datasets with their status, metrics, and health information on a consolidated list.
2. Click on a dataset to view its detailed configuration and insights.
3. Retire unused datasets to maintain an organized environment.

<figure><img src="/files/LIoS0MRmG1K4zlutlKZz" alt=""><figcaption><p>Live Dataset List</p></figcaption></figure>

<figure><img src="/files/XocvSlPWIIY3xuE27YxB" alt=""><figcaption><p>Retired Dataset List</p></figcaption></figure>

***

## Alerts and Notifications

1. Create custom alerts to monitor specific events related to datasets and connectors.
2. Set up notifications for events, ensuring timely action when needed.
3. Configure alerts to notify you via channels like Slack, email, and Discord for real-time updates.

<figure><img src="/files/5VaEO7Kdmb1Iz0VNYSCa" alt=""><figcaption><p>Infra Alerts</p></figcaption></figure>

<figure><img src="/files/x5t9MJYv88JS91mh21y2" alt=""><figcaption><p>Custom Alerts</p></figcaption></figure>

<figure><img src="/files/V5Ay5jEq9xpgb06VVCdl" alt=""><figcaption><p>Notification Channels</p></figcaption></figure>

***

## Dashboards

1. Monitor the overall system health with infra cluster and health dashboards.
2. Gain insights into ingestion, processing, storage, and querying workflows through dedicated dashboards.

<figure><img src="/files/ZnxfPdMnmKb3W9K6rG7D" alt=""><figcaption><p>Infra Metrics</p></figcaption></figure>

<figure><img src="/files/M18uUaFwtwG2Uo7MCjZv" alt=""><figcaption><p>Ingestion Metrics</p></figcaption></figure>

<figure><img src="/files/eh95q6ByJrJBMoH34hZZ" alt=""><figcaption><p>Storage Metrics</p></figcaption></figure>

<figure><img src="/files/6KfgR6OnQvCw46Pza8ls" alt=""><figcaption><p>Processing Metrics</p></figcaption></figure>

***

The **Management Console** empowers you to manage your datasets, ensure data integrity, monitor system performance, and stay proactive with alerts and notifications—all in one streamlined interface.


# Developer Guide

Guide to setup and run in a developer machine

### Introduction: <a href="#id-5zntod7zxfl6" id="id-5zntod7zxfl6"></a>

In this document, we'll explore how to fork, clone, build, and install Obsrv from these GitHub repositories. Obsrv uses GitHub repositories as they are versatile tools for managing and sharing code, documentation, and other resources efficiently.

### Prerequisites: <a href="#lwdfmye9jg68" id="lwdfmye9jg68"></a>

Before we begin, ensure you have the following:

* A GitHub account
* Basic knowledge of Git and GitHub concepts.
  * Fork the Repository - [reference doc](https://docs.github.com/en/pull-requests/collaborating-with-pull-requests/working-with-forks/fork-a-repo)
  * Clone the Repository - [reference doc](https://docs.github.com/en/repositories/creating-and-managing-repositories/cloning-a-repository)
* IDE for development. Suggested - Intellij idea. Installation reference doc can be found under [IntelliJ](https://www.jetbrains.com/help/idea/installation-guide.html)
  * Open the Project - [reference doc](https://www.jetbrains.com/guide/java/tutorials/import-project/open-project/)
  * Build & Run
* Docker account & Docker desktop application - [reference doc](https://docs.docker.com/desktop/)
* Cloud provider account. Currently Obsrv supports - AWS, AZURE & GCP.
* Specific details about each repository are listed under the repository section.

### Obsrv GitHub repositories: <a href="#at5woe3fkmsh" id="at5woe3fkmsh"></a>

#### **obsrv-core:**

**Repository:** <https://github.com/Sunbird-Obsrv/obsrv-core>

**Introduction:**

Obsrv-core is a framework consisting of Flink jobs designed to handle data extraction and processing tasks efficiently. It provides a flexible and customizable pipeline for various data-related operations. These jobs have been designed to process, enrich, and validate data from various sources, making them highly adaptable to a wide range of datasets. The data streaming jobs are built with a generic approach that makes them robust and able to handle diverse datasets without requiring significant changes to the underlying code. More details on Obsrv-core can be found here - <https://github.com/Sunbird-Obsrv/obsrv-core/blob/main/README.md>

**Prerequisite:**

To use the Obsrv core, make sure you have the following dependencies installed:

* Java 11
* Maven
* Docker

**Build Project:**

* Use the cd command to navigate to the root directory of obsrv-core project
* Run the following Maven command to build the project - ​​mvn clean install

**Build Artifact:**

* Run the following Docker command from the root directory of obsrv-core to build the artifact - ​​docker build .
* Tag the built image with name - docker tag \<digest-from-build-command> \<image-name>:\<version>

Sample command: docker tag sha256:6021feb600a9ec4139dadd06536f46e5201435f0837fd4f60581d41a6ae6cf63 sunbird/obsrv-core:1.0.0

* Push docker image to docker hub - docker push \<image-name>:\<version>

#### **obsrv-api-service:**

**Repository:** <https://github.com/Sunbird-Obsrv/obsrv-api-service>

**Introduction:**

This repository mainly has 2 sets of APIs. Obsrv-api-service and Obsrv-command-service.

**Obsrv-api-service** is a set of APIs that provide access to a variety of data sources and datasets. These APIs can be used to query and analyze different types of events, as well as to manage data sources and datasets. More details on Obsrv-api-service can be found here - <https://github.com/Sunbird-Obsrv/obsrv-api-service/blob/main/README.md>

**Obsrv-command-service** is an API that provides access to restart the Obsrv pipeline to make a newly created dataset accessible on the Obsrv.

**Prerequisite:**

To use the Obsrv API service, make sure you have the following dependencies installed:

* Node.js: version 18
* TypeScript: version 4.8.4
* Express.js: version 4.18.2
* npm: version 9.6.4
* Docker

To use the Obsrv Command service, make sure you have the following dependencies installed:

* Python3
* pip3
* Docker

**Build Artifact:**

Obsrv API service:

* Use the following command to navigate to the api-service directory of obsrv-api-service project - cd obsrv-api-service/api-service
* Install the required dependencies by running the following command - npm install

Obsrv Command service:

* Use the following command to navigate to the command-service directory of obsrv-command-service project - cd obsrv-api-service/command-service
* Install the required dependencies by running the following command - pip install -r requirements.txt

Artifact build commands:

* Use the following command to navigate to the respective service directory of obsrv-api-service project - cd obsrv-api-service/api-service or cd obsrv-api-service/command-service
* Run the following Docker command from the respective services directory to build the artifact - ​​docker build .
* Tag the built image with name - docker tag \<digest-from-build-command> \<image-name>:\<version>

Sample command: docker tag sha256:8021feb600a9ec4139dadd06536f48e5201435f0837fd4f60581d41a6ae6cf83 sunbird/obsrv-api-service:1.0.0

* Push docker image to docker hub - docker push \<image-name>:\<version>

#### **obsrv-automation:**

**Repository:** <https://github.com/Sunbird-Obsrv/obsrv-automation>

**Introduction:**

Obsrv-automation repository provides support for installation of Obsrv across major cloud providers. More details of Obsrv installation on different cloud providers can be found here - <https://github.com/Sunbird-Obsrv/obsrv-automation/blob/main/README.md>

**Prerequisite:**

* Terragrunt: 0.45.6 - Please see [Install Terragrunt](https://terragrunt.gruntwork.io/docs/getting-started/install/) for reference.
* Terraform: 1.5.7
* Terrahelp: 0.7.5
* Kubectl
* Helm: 3.10.2
* Azure cli or AWS cli: 2.13.8

**Steps to install Obsrv:**

* Setup prerequisites depending on the choice of cloud provider from this [doc](https://github.com/Sunbird-Obsrv/obsrv-automation/blob/main/README.md)
* Update the automation scripts to use the new docker images created if any from other repositories under terraform/\<cloud-provider>/variables.tf file.
* Below are the installation steps to setup complete Obsrv

```
cd terraform/<cloud-provider>
terragrunt init
terragrunt plan
terragrunt apply
```


# Example Datasets

Download the Postman Collection from [here](https://github.com/Sanketika-Obsrv/obsrv-api-service/blob/main/api-service/postman-collection/Obsrv%20v2%20apis.postman_collection.json)

The above collection has the complete list of API's, which include

1. Creation of Dataset
2. Publishing of Datasets
3. Ingestion of Sample Events
4. Querying of Ingested Events using both SQL and Native API's

You are free to play around, once the installation is completed successfully


# OpenTelemetry (OTEL) Integration

### Overview

[OpenTelemetry (OTEL)](https://opentelemetry.io/) is a vendor-neutral, open standard for collecting observability data from software systems. It provides a unified way to capture three types of signals:

* **Traces (Spans)** — track the execution of individual operations across services, capturing timing, status, and contextual attributes for each step in a request lifecycle.
* **Logs** — record discrete events that occurred within a service, such as audit actions, state changes, errors, or informational messages, along with structured metadata.
* **Metrics** — report aggregated measurements over time, such as latency percentiles, success rates, and error counts, grouped by granularity and frequency.

#### Why OTEL with Obsrv?

Raw OTEL payloads use a deeply nested, key-value array structure (e.g., `[{"key": "foo", "value": {"stringValue": "bar"}}]`) that is designed for portability but is not directly queryable. Obsrv addresses this by providing an **OTEL service** that acts as a translation layer — it receives raw OTEL signals, flattens them into clean JSON objects, and feeds both representations into the Obsrv data pipeline via Kafka.

This means you can instrument your services with standard OTEL SDKs without any Obsrv-specific changes, and still get fully indexed, queryable datasets with no manual schema wrangling.

#### Supported signal types

| Signal  | OTEL payload key  | Obsrv `eid` | Use case                                           |
| ------- | ----------------- | ----------- | -------------------------------------------------- |
| Traces  | `resourceSpans`   | `API`       | API call tracing, latency, error tracking          |
| Logs    | `resourceLogs`    | `AUDIT`     | Audit trails, state change events, error logs      |
| Metrics | `resourceMetrics` | `METRIC`    | Aggregated KPIs — latency, throughput, error rates |

#### Data flow

<figure><img src="/files/D5iVw0b0GUlgQKbpY1Ic" alt=""><figcaption></figcaption></figure>

**What the OTEL service does:**

1. Receives raw OTEL payloads over HTTP.
2. Unpacks the nested attribute arrays and `kvlistValue` structures into flat key-value JSON.
3. Promotes resource-level fields (`eid`, `producer`, `producerType`) and scope-level fields (`name`, `version`, `scope_uuid`, `count`) to the top level.
4. Writes the **raw** event and the **flattened** event to two separate Kafka topics in parallel.
5. The flattened events are consumed by Obsrv and indexed into the configured dataset.

***

### Step 1 — Emit OTEL events from your service

Use any OpenTelemetry SDK to instrument your service. The OTEL service accepts all three standard signal types.

#### Traces (Spans)

Emitted as `resourceSpans`. Each span captures a single operation — e.g., an API request — including timing, status, and any error events.

```json
{
  "resourceSpans": [{
    "resource": {
      "attributes": [
        { "key": "eid",           "value": { "stringValue": "API" } },
        { "key": "producer",      "value": { "stringValue": "iam" } },
        { "key": "producer.type", "value": { "stringValue": "Facilitator" } }
      ]
    },
    "scopeSpans": [{
      "scope": {
        "name": "iam_service",
        "version": "1.0.0",
        "attributes": [
          { "key": "scope.uuid", "value": { "stringValue": "874d0df2-224a-4bea-9e0c-303995a38937" } },
          { "key": "count",      "value": { "intValue": 1 } }
        ]
      },
      "spans": [{
        "name": "userTokenGenerationV1",
        "startTimeUnixNano": "1747805502725834184",
        "endTimeUnixNano":   "1747805507924834184",
        "status": "Ok",
        "traceId": "a923d739-774a-4e90-959c-21f064e586f1",
        "spanId":  "95cce13d-bfb1-4ce6-9891-ad7d2d21ac03",
        "attributes": [
          { "key": "sender.id",            "value": { "stringValue": "ma****@test.in" } },
          { "key": "span.uuid",            "value": { "stringValue": "a1e082c3-3c36-40cb-bec5-643925259a2f" } },
          { "key": "observed.time.unix.nano", "value": { "stringValue": "1747805502725806984" } },
          { "key": "request.method",       "value": { "stringValue": "POST" } },
          { "key": "request.url",          "value": { "stringValue": "/v1/user/token/generate" } },
          { "key": "response.status.code", "value": { "stringValue": "200" } },
          { "key": "response.time.ms",     "value": { "stringValue": "5199" } }
        ]
      }]
    }]
  }]
}
```

When a span contains an error, include an `events` array:

```json
"events": [{
  "name": "error",
  "time": "2025-05-21T05:33:48.324173Z",
  "attributes": [
    { "key": "msg",  "value": { "stringValue": "Invalid username/password" } },
    { "key": "code", "value": { "stringValue": "TKN_2025" } },
    { "key": "type", "value": { "stringValue": "UNAUTHORIZED" } }
  ]
}]
```

#### Logs

Emitted as `resourceLogs`. Each log record captures a discrete event — e.g., an audit action — with severity, body text, and structured state attributes.

```json
{
  "resourceLogs": [{
    "resource": {
      "attributes": [
        { "key": "eid",          "value": { "stringValue": "AUDIT" } },
        { "key": "producer",     "value": { "stringValue": "iam-service" } },
        { "key": "producerType", "value": { "stringValue": "IAM" } }
      ]
    },
    "scopeLogs": [{
      "scope": {
        "name": "iam_service",
        "version": "1.0.0",
        "attributes": [
          { "key": "scope_uuid", "value": { "stringValue": "df928931-107e-4539-9b21-e01f3a137b53" } },
          { "key": "count",      "value": { "intValue": 1 } }
        ]
      },
      "logRecords": [{
        "timeUnixNano":         1747138337071,
        "observedTimeUnixNano": "1747138337071729909",
        "severityNumber": "12",
        "traceId": "ee36cf84-ce43-42a3-8384-a6b087fc0823",
        "spanId":  "2dc6e05c-68ac-4b67-880e-4afd85784989",
        "body": { "stringValue": "User Create" },
        "attributes": [
          { "key": "type",     "value": { "stringValue": "User" } },
          { "key": "status",   "value": { "stringValue": "OK" } },
          { "key": "id",       "value": { "stringValue": "te**********@yopmail.com" } },
          { "key": "log_uuid", "value": { "stringValue": "16ce4388-4ce5-4921-9b2e-b024aec28f62" } },
          { "key": "state", "value": { "kvlistValue": { "values": [
            { "key": "email",     "value": { "stringValue": "te**********@yopmail.com" } },
            { "key": "username",  "value": { "stringValue": "te**********@yopmail.com" } },
            { "key": "firstname", "value": { "stringValue": "test" } },
            { "key": "lastname",  "value": { "stringValue": "user" } }
          ]}}}
        ]
      }]
    }]
  }]
}
```

#### Metrics

Emitted as `resourceMetrics`. Each metric payload can carry multiple named measurements — e.g., latency percentiles, error rates — bundled under a single scope and aggregated over a time window.

```json
{
  "resourceMetrics": [{
    "resource": {
      "attributes": [
        { "key": "eid",          "value": { "stringValue": "METRIC" } },
        { "key": "producer",     "value": { "stringValue": "APP1" } },
        { "key": "producerType", "value": { "stringValue": "App" } }
      ]
    },
    "scopeMetrics": [{
      "scope": {
        "name": "metrics_service",
        "version": "1.0",
        "attributes": [
          { "key": "scope_uuid", "value": { "stringValue": "9db4-325096b39f47" } },
          { "key": "checksum",   "value": { "stringValue": "120EA8A25E5D487BF68B5F7096440019" } },
          { "key": "count",      "value": { "intValue": 8 } }
        ]
      },
      "metrics": [
        {
          "name": "latency_avg_ms",
          "unit": "ms",
          "sum": {
            "aggregationTemporality": 1,
            "isMonotonic": false,
            "dataPoints": [{
              "asDouble": 1153.1718,
              "startTimeUnixNano": "1544712660000000000",
              "endTimeUnixNano":   "1544712661590000000",
              "attributes": [
                { "key": "metric_uuid",        "value": { "stringValue": "43kr3d5f-3cfb-4e6e-b6a2-0ee5d6508923" } },
                { "key": "observedTimeUnixNano","value": { "stringValue": "1581452772000000321" } },
                { "key": "metric.code",        "value": { "stringValue": "latency_avg_ms" } },
                { "key": "metric.category",    "value": { "stringValue": "Usage" } },
                { "key": "metric.label",       "value": { "stringValue": "Average Latency in ms" } },
                { "key": "metric.granularity", "value": { "stringValue": "minute" } },
                { "key": "metric.frequency",   "value": { "stringValue": "10-min" } }
              ]
            }]
          }
        }
      ]
    }]
  }]
}
```

> **Note:** Multiple metrics can be included in a single `scopeMetrics.metrics` array. The sample data includes `latency_avg_ms`, `latencyP50_ms`, `latencyP95_ms`, `latencyP99_ms`, `success_percent`, `timeout_percent`, `server_error_percent`, and `client_error_percent`.

> **Timing for metrics:** `startTimeUnixNano` should be the window start (e.g., 10 minutes before the request), and `endTimeUnixNano` should be the time of the request.

***

### Step 2 — Set up and run the OTEL service

The OTEL service is the transformation layer between your instrumented services and Obsrv. It receives raw OTEL payloads, flattens them into a queryable structure, and publishes both representations to Kafka.

**Repository:** <https://github.com/Sanketika-Obsrv/otel-service>

***

#### Prerequisites

* **Node.js** 20+
* **Kafka** broker accessible from the service
* An Obsrv dataset already created (or planned) — the `dataset-id` in the API path must match the Obsrv dataset ID

***

#### Installation and setup

**Clone the repository:**

```bash
git clone https://github.com/Sanketika-Obsrv/otel-service.git
cd otel-service
```

**Install dependencies and start locally:**

```bash
npm install
npm start
```

The service starts on port `3000` by default.

**Run with Docker:**

```bash
docker build -t otel-service:latest .
docker run -p 3000:3000 \
  -e kafka_host=<your-kafka-host> \
  -e kafka_port=9092 \
  -e system_env=dev \
  otel-service:latest
```

**Deploy to Kubernetes using Helm:**

```bash
helm install otel-service ./helm-chart \
  --set config.kafka_host=<your-kafka-host> \
  --set SYSTEM_ENV=dev
```

The Helm chart deploys the service as a `LoadBalancer` in the `otel-api` namespace with a Prometheus `ServiceMonitor` included.

***

#### Configuration

All configuration is controlled via environment variables:

| Variable           | Default              | Description                                                   |
| ------------------ | -------------------- | ------------------------------------------------------------- |
| `port`             | `3000`               | Port the HTTP server listens on                               |
| `kafka_host`       | `localhost`          | Kafka broker hostname                                         |
| `kafka_port`       | `9092`               | Kafka broker port                                             |
| `system_env`       | `local`              | Environment prefix for Kafka topic names (e.g. `dev`, `prod`) |
| `ingest_topic`     | `ingest`             | Topic name suffix for flattened/transformed events            |
| `otelingest_topic` | `otelingest`         | Topic name suffix for raw OTEL events                         |
| `app_name`         | `obsrv-otel-service` | Service name used in logs and metrics                         |

**Kafka topic naming convention:**

Topics are constructed as `{system_env}.{topic_suffix}`. With defaults:

| Topic            | Default name       | Contents                                   |
| ---------------- | ------------------ | ------------------------------------------ |
| Flattened events | `local.ingest`     | Transformed, flat JSON — consumed by Obsrv |
| Raw OTEL events  | `local.otelingest` | Original nested OTEL payload               |

For a `dev` environment with custom topic names:

```
dev.ingest       ← set system_env=dev, ingest_topic=ingest
dev.otelingest   ← set system_env=dev, otelingest_topic=otelingest
```

***

#### Sending events to the OTEL service

**Endpoint:**

```
POST http://localhost:3000/network-observability/v1/in/<dataset-id>
Content-Type: application/json
```

| Parameter    | Description                                                 |
| ------------ | ----------------------------------------------------------- |
| `dataset-id` | The Obsrv dataset ID created in Step 3. Must match exactly. |

**Request body format:**

Wrap your raw OTEL payload inside a `data` object:

```json
{
  "data": {
    "resourceSpans": [ ... ]
  }
}
```

The same wrapper applies for logs and metrics:

```json
{
  "data": {
    "resourceLogs": [ ... ]
  }
}
```

```json
{
  "data": {
    "resourceMetrics": [ ... ]
  }
}
```

> The service detects the signal type automatically — if the payload contains `resourceSpans`, `resourceLogs`, or `resourceMetrics` it is treated as a standard OTEL v2 event and transformed accordingly.

**Successful response (HTTP 200):**

```json
{
  "id": "otel.data.in",
  "ver": "v1",
  "ets": 1747393761822,
  "params": {
    "resmsgid": "792c99b3-ab06-459e-86a3-2da6902306c6",
    "err": "",
    "status": "SUCCESSFUL",
    "errmsg": ""
  },
  "responseCode": "OK",
  "result": {
    "message": "The data has been successfully ingested"
  }
}
```

**Error response (HTTP 500):**

```json
{
  "id": "otel.data.in",
  "ver": "v1",
  "ets": 1747393761822,
  "params": {
    "resmsgid": "...",
    "err": "SERVER_ERROR",
    "status": "FAILED",
    "errmsg": "Error message"
  },
  "responseCode": "SERVER_ERROR",
  "result": {}
}
```

**Maximum request body size:** 5 MB per request.

***

#### What the service does internally

For each ingested event, the OTEL service:

1. Validates the request body schema (requires a `data` object).
2. Attaches the `datasetId` (from the URL path) to the payload.
3. Detects the event format — standard OTEL (`resourceSpans` / `resourceLogs` / `resourceMetrics`) or legacy v1.
4. **Transforms the payload:** unpacks OTel attribute arrays (`[{"key": ..., "value": ...}]`) and `kvlistValue` structures into flat key-value maps, promotes resource and scope fields to the top level, and explodes metric `dataPoints` into individual events.
5. Adds processing metadata: `mid` (message ID), `syncts` (sync timestamp), and `obsrv_meta` (source and routing info).
6. Publishes the **flattened event** to `{system_env}.{ingest_topic}` and the **raw event** to `{system_env}.{otelingest_topic}` in parallel.
7. Increments Prometheus counters for monitoring.

***

### Step 3 — Create an Obsrv dataset using the flattened schema

After the OTEL service transforms events, the flattened output looks like this:

#### Flattened Span (API trace)

```json
{
  "resource": {
    "eid": "API",
    "producer": "iam",
    "producer.type": "Facilitator"
  },
  "scope": {
    "name": "iam_service",
    "version": "1.0.0",
    "attributes": {
      "scope.uuid": "0dae9978-51cb-496d-882d-c6d632e52cba",
      "count": 1
    }
  },
  "edata": {
    "name": "userTokenGenerationV1",
    "startTimeUnixNano": "1747393760692953530",
    "endTimeUnixNano": "1747393761278953530",
    "status": "UNAUTHORIZED",
    "traceId": "fd8c2cea-f34a-43df-8c9c-dea27f338960",
    "spanId": "479c8550-fda2-4ba8-9446-2cf140b5042b",
    "mid": "792c99b3-ab06-459e-86a3-2da6902306c6",
    "ets": 1747393761822,
    "attributes": {
      "sender.id": "ma****@test.in",
      "span.uuid": "792c99b3-ab06-459e-86a3-2da6902306c6",
      "observed.time.unix.nano": "1747393760692942330",
      "request.method": "POST",
      "request.url": "/v1/user/token/generate",
      "response.status.code": "401",
      "response.time.ms": "586"
    },
    "events": {
      "error": {
        "time": "2025-05-16T11:09:20.692589Z",
        "attributes": {
          "msg": "Invalid username/password",
          "code": "TKN_2025",
          "type": "UNAUTHORIZED"
        }
      }
    }
  }
}
```

#### Flattened Log (Audit)

```json
{
  "resource": {
    "eid": "AUDIT",
    "producer": "iam-service",
    "producerType": "IAM"
  },
  "scope": {
    "name": "iam_service",
    "version": "1.0.0",
    "attributes": {
      "scope_uuid": "ab9ba7fa-3954-46aa-b176-b7cc53671b56",
      "count": 1
    }
  },
  "edata": {
    "timeUnixNano": 1747138417347,
    "observedTimeUnixNano": "1747138417347597974",
    "severityNumber": "12",
    "traceId": "f8ddc431-5ce5-4c08-8adb-15e3ba7ee2c3",
    "spanId": "6bf31af5-df0f-4d3e-ab3c-e0c6ceda449c",
    "body": "User already exists with email testemail13k@yopmail.com",
    "mid": "695c032a-4a6a-4c59-a8fc-936ebb6bbd1d",
    "ets": 1747138418315,
    "attributes": {
      "state": {
        "email": "te**********@yopmail.com",
        "username": "te**********@yopmail.com",
        "firstname": "test",
        "lastname": "user"
      },
      "type": "User",
      "log_uuid": "695c032a-4a6a-4c59-a8fc-936ebb6bbd1d",
      "status": "Failed",
      "id": "te**********@yopmail.com"
    }
  }
}
```

Use this flattened structure to define your dataset schema in Obsrv.

***

### Step 4 — Query your data

Once the dataset is active and events are flowing through the pipeline:

* **Superset** — Use the Obsrv-connected Superset instance to build charts and dashboards directly on the dataset.
* **Data Out API** — Query programmatically via the Obsrv Data Out API.

***


# How Tos

Step-by-step recipes for common Obsrv customizations

This section contains step-by-step recipes for customizing Obsrv beyond the standard console workflows. Each how-to lists the exact configuration changes, SQL statements and deployment steps required, along with the trade-offs of each approach.

## Available How Tos

* [Extending Obsrv with Custom Processing Logic](/guides/how-tos/custom-code-injection) — insert your own code (Python/Node.js or any Kafka client) alongside the dataset processing flow.


# Extending Obsrv with Custom Processing Logic

Insert your own code (Python, Node.js or any Kafka consumer/producer) alongside the Obsrv dataset processing flow.

## Why?

Obsrv datasets already support JSONata expressions for transformation — for straightforward field mapping and reshaping, that's the standard path, no infra change needed. But JSONata has limits: it can't call external systems, run extensive/multi-step custom logic, or do anything beyond expression evaluation.

When your transformation genuinely needs more than JSONata can express — before the data lands in storage, whether that's analytics (Druid) or a transactional/lakehouse store — you need your own code running as a step in the pipeline. That's what a custom job is: your own Kafka consumer/producer (Python, Node.js, or any language) inserted into the event flow. It's **not tied to any particular purpose** — it can enrich events, filter them, call external systems, reshape fields, or something else entirely.

## How?

There are three solutions:

1. **Solution 1: Drop the custom stream job in between the unified pipeline** — switch the unified pipeline to individual jobs, then insert custom logic between any two stages. Covered in "Solution 1" below.
2. **Solution 2: After the pipeline router** — keep the unified pipeline as-is, and insert custom logic after the router by fanning out a second consumer on its output topic. Covered in "Solution 2".
3. **Solution 3: Write a custom connector** — run your transformation logic inside the source connector itself, before the event ever reaches the pipeline. Covered in "Solution 3".

| Your situation                                                                                    | Use                                                |
| ------------------------------------------------------------------------------------------------- | -------------------------------------------------- |
| Custom code between two specific stages (needs pipeline-enriched input, e.g. denormalized fields) | Solution 1: individual jobs + processor in the gap |
| Custom code on final routed events, and comfortable owning a manual Druid cutover                 | Solution 2: After the pipeline router              |
| Logic can run entirely at the source, before any event reaches Obsrv                              | Solution 3: Write a custom connector               |

## How the pipeline is structured today

Every pipeline stage is an independent Flink job with a Kafka **in** topic and a Kafka **out** topic:

| Job            | In topic    | Out topic                                                                                      | Failed topic       |
| -------------- | ----------- | ---------------------------------------------------------------------------------------------- | ------------------ |
| extractor      | `ingest`    | `raw`                                                                                          | `failed`           |
| preprocessor   | `raw`       | `unique`                                                                                       | `failed`           |
| denormalizer   | `unique`    | `denorm`                                                                                       | `failed`           |
| transformer    | `denorm`    | `transform`                                                                                    | `transform.failed` |
| dataset-router | `transform` | per-dataset `router_config.topic` (the topic this dataset's events land in — e.g. `d1-events`) | `failed`           |

**Existing Unified Pipeline flow:**

<figure><img src="/files/G6wCePRJ5QsNvHPXNrsC" alt="Existing pipeline flow: producers and connectors through ingest, extractor, preprocessor, denormalizer, transformer, dataset-router, to Druid"><figcaption></figcaption></figure>

## If your job changes the event shape

If your custom job adds a new field, or changes the data type of a field that already exists, update the dataset's schema in the console **before** that job's output reaches Druid. Skip this and a new field is silently dropped, or a type-changed field fails ingestion outright — since Druid ingests strictly against the schema it already knows.

1. Open the dataset in the Obsrv management console.
2. Go to **Schema Details** → **Ingestion**.
3. New field your job adds (e.g. `meta`): click **Add Field**, enter the exact field name your job writes, and pick its type.
4. Existing field whose type your job changes: find it in the list and update its **Type** to match what your job now emits.
5. Save and publish — this regenerates Druid's ingestion spec with the updated field list before any reshaped event arrives.

You (or whoever writes the custom job) already know what it does to the event, so make this schema change first, then deploy the job.

## Solution 1: Drop the custom stream job in between the unified pipeline

This requires switching from the unified pipeline to individual jobs first — five separate Flink jobs instead of one. Instead of deploying the unified pipeline from the automation charts, disable it and deploy each job individually — refer to [obsrv-core](https://github.com/Sanketika-Obsrv/obsrv-core) for each job's build and deployment steps.

Because every stage reads from a topic and writes to a topic, inserting custom code between **any** two stages is always the same three moves:

1. **Rename the downstream stage's in topic** to a new "pre" topic (one line in that job's configuration).
2. **Run your custom job** consuming the upstream stage's **unchanged, stock out topic**, and producing to the new "pre" topic.
3. **Leave every other job untouched — including the upstream stage.**

<figure><img src="/files/Sbrhu4FNy6FL1wVqRjez" alt="General insertion pattern: before shows Job 1 producing to Topic 1, consumed by Job 2; after shows Job 1 still producing to the unchanged Topic 1, consumed by the custom streaming job, which produces to a renamed Topic 1_pre, consumed by Job 2"><figcaption></figcaption></figure>

For example, inserting between **denormalizer and transformer**:

1. Override the downstream job's in topic in its configuration:

```hocon
# transformer job config
kafka {
  input.topic = "transform_pre"   # stock value: "denorm"
  output.transform.topic = "transform"
}
```

2. Create the `transform_pre` topic (partition count = the stock topic's partition count).
3. Run the custom job with `IN_TOPIC=denorm`, `OUT_TOPIC=transform_pre`.
4. Denormalizer (`unique` → `denorm`) and router (`transform` → dataset topics) stay on stock configuration.
5. If this job adds or retypes fields, do the schema update in "If your job changes the event shape" above first.

Resulting flow:

<figure><img src="/files/JSPa36vXUTgAl52IbmLY" alt="Custom job inserted between denormalizer and transformer: denormalizer produces to denorm unchanged, the custom job consumes it and produces to transform_pre, transformer&#x27;s in topic is repointed to transform_pre while dataset-router continues unchanged through to Druid"><figcaption></figcaption></figure>

The upstream stage never knows the difference — it keeps producing to its stock out topic; events simply pass through your code first. For any other gap, substitute the topic pair. E.g. between preprocessor and denormalizer: denormalizer `input.topic = "unique_pre"`, processor `unique` → `unique_pre`. The in-topic key per job: extractor `kafka.input.topic`, preprocessor `input.topic`, denormalizer `input.topic`, transformer `input.topic`, dataset-router `input.topic`.

### Pick the insertion point and rewire one topic

Pick the gap based on **what your code needs as input** — and what shape the events are in at that point:

| Insertion point             | Events carry                                                                                     | Typical use                        |
| --------------------------- | ------------------------------------------------------------------------------------------------ | ---------------------------------- |
| before extractor (`ingest`) | batch envelope: `{"dataset": "...", "events": [...]}` (or `{"event": {...}}` for a single event) | normalize source payloads          |
| preprocessor → denormalizer | single-event envelope: `{"event": {...}, "obsrv_meta": {...}}`                                   | enrich before denorm lookups       |
| denormalizer → transformer  | single-event envelope, with denormalized fields present                                          | logic that needs master-data joins |
| transformer → router        | single-event envelope, with JSONata outputs present                                              | post-process transformed fields    |

When unwrapping a single-event envelope, modify the nested `event` object and pass `obsrv_meta` through unchanged — it carries stage flags and timings the rest of the pipeline relies on.

## Solution 2: After the pipeline router

1. If your job adds or retypes fields, do the schema update in "If your job changes the event shape" above first, and publish.
2. Deploy your custom job: consume the dataset's live topic (e.g. `user-data`) as a second consumer group, produce to a new topic with a different name (e.g. `user-data-processed`). See "Reference: the custom streaming job" below for a minimal example.
3. Confirm the job is healthy and producing correctly-shaped events to the new topic.
4. Take the existing supervisor's spec, change `dataSchema.dataSource` to a new name and point `ioConfig.topic` at the new topic, and submit it via the Druid console or Supervisor API — this creates a new, separate Druid datasource ingesting the processed data.
5. Once the new datasource's supervisor is healthy and ingesting correctly, suspend the original, obsrv-created datasource's supervisor.

Queries/dashboards on this dataset now need to point at the new datasource name.

## Solution 3: Write a custom connector

1. Base your connector on an existing open-source Obsrv connector — e.g. [jdbc-connector](https://github.com/Sanketika-Obsrv/jdbc-connector) — and adapt it for your source/use case.
2. Run your transformation logic inside the connector itself, before it produces the event — see the [connectors developer guide](/guides/connectors-developer-guide) for interfaces and packaging.
3. Your connector produces wherever the reference connector already produces to — no topic rewiring or Druid cutover needed.
4. If your connector adds or retypes fields, still do the schema update in "If your job changes the event shape" above first.
5. Package and deploy per the connector guide's packaging steps.

## Reference: the custom streaming job

For reference, not a required step — any Kafka client works for the custom streaming job used above, it just needs to consume `IN_TOPIC`, run your logic, and produce to `OUT_TOPIC`. Matching Solution 2, where events are flat, in Python:

```python
for msg in consumer:                       # consume IN_TOPIC
    event = json.loads(msg.value())        # after-router output is flat — no wrapper
    event["metadata"] = my_custom_logic(event)   # <- your code
    producer.produce(OUT_TOPIC, json.dumps(event).encode())
```

That's it — the pipeline stages on either side don't need to know it's there. If you're inserting **between individual jobs** instead (Solution 1), events are wrapped in an envelope, not flat — see "Pick the insertion point and rewire one topic" above for the exact shape to unwrap.

{% hint style="info" %}
Obsrv Flink producers use **snappy** compression, so `confluent-kafka` (Python) works out of the box; Node.js `kafkajs` needs the `kafkajs-snappy` codec registered.
{% endhint %}


# Connectors Developer Guide

Connectors, alongside datasets, are fundamental components of Obsrv. While datasets establish a framework for managing data processing, storage, and querying, connectors provide a framework for managing data ingress and egress.

This section will guide you through the process of developing, packaging, and deploying connectors using the provided SDKs and repositories. The connectors can be written in Scala, Java, or Python and can work with various data processing frameworks like Spark and Flink.

## Connector Types

Batch connectors are designed to handle data in chunks or sets, enabling efficient processing of large volumes of data that do not require immediate handling. These connectors are ideal for use cases where data can be accumulated and processed at scheduled intervals. In contrast, streaming connectors are built to process data in real-time, providing instant insights and immediate data handling as it flows into the system. Streaming connectors are essential for applications that demand continuous data updates and low-latency processing, such as real-time analytics and monitoring systems. Together, these connectors facilitate seamless data integration across varying use cases and workloads.

### Source Connectors

#### Batch Connectors

Batch connectors are commonly used to pull data from various sources such as databases and file systems. These connectors operate over accumulated data over time and then processing it collectively at scheduled intervals. This approach is especially effective when working with data stored in database systems or file systems by using frameworks such as Apache Spark, which can handle large datasets efficiently. By executing these tasks on a scheduled basis, batch connectors optimize resource usage and ensure that large volumes of data are processed with minimal latency.

{% hint style="info" %}
Current versions of SDKs are available for Java, Scala, and Python.
{% endhint %}

#### Streaming Connectors

Streaming connectors that run on Apache Flink are crucial for integrating and processing real-time data streams from platforms like Apache Kafka, Amazon Kinesis, and Debezium. These connectors facilitate the seamless ingestion and processing of data in a distributed, fault-tolerant manner, enabling applications to react to data as it arrives. By leveraging Flink's robust stateful processing capabilities, these connectors support complex event processing, real-time analytics, and continuous data transformations, making them ideal for scenarios requiring immediate insights and low-latency responses.

{% hint style="info" %}
Current versions of SDKs are available for Java and Scala.
{% endhint %}

### Sink Connectors

Sink connectors are used to export data from centralized data platforms, like Obsrv, to various external systems and storage solutions. These connectors ensure that processed data can be efficiently moved to destinations such as data warehouses, cloud storage services, and analytics platforms. By leveraging sink connectors, businesses can maintain updated datasets across multiple systems, enabling actionable insights and enhanced data-driven decision-making.

{% hint style="info" %}
Current versions of SDKs don’t support Sink Connectors and are under development and not available as of now.
{% endhint %}


# SDK Assumptions

## Package Strucuture

For the connector to work seamlessly with Obsrv, the connector package has to be in the following structure.

{% tabs %}
{% tab title="Java/Scala" %}

```
example-connector-0.1.0-distribution.tar.gz
├── libs
    └── sample-dependency.jar
├── example_connector.jar
├── alerts.yaml
├── metadata.json
├── metrics.yaml
├── ui-config.json
└── icon.svg
```

{% endtab %}

{% tab title="Python" %}

```
example-connector-0.1.0-distribution.tar.gz
├── libs # optional
    └── sample-dependency.jar
├── example_connector
    └── main.py
├── alerts.yaml
├── metadata.json
├── metrics.yaml
├── requirements.txt
├── ui-config.json
└── icon.svg
```

> NOTE: The `libs` in python distribution are only required, if PySpark has any dependent JARs.
> {% endtab %}
> {% endtabs %}

## Supported Event Types

### Stream Connectors

1. When emitting an event to the SDK, the SDK only supports JSON format. If the format is other than JSON, events will be skipped with an error of type `INVALID_DATA_FORMAT_ERROR`
2. In Stream Connectors, the ID has to be unique when registering a stream, if its a CONSTANT. If the ID is being referenced from the Configuration, no action has to be taken.

### Batch Connectors

The batch connector expects a [Spark DataFrame](https://spark.apache.org/docs/latest/sql-programming-guide.html#datasets-and-dataframes) when its returned to the SDK


# Required Files

The SDK expects the following files to be as part of the source distribution:

* [**metadata.json**](/guides/connectors-developer-guide/required-files/metadata-json): This file is used mainly for registering a connector with the Obsrv system.
* [**ui-config.json**](/guides/connectors-developer-guide/required-files/ui-config-json): This JSON configuration schema is used to define the data that the user interface (UI) will collect from the user to configure the connector.
* [**metrics.yaml**](/guides/connectors-developer-guide/required-files/metrics.yaml): This YAML configuration is used to define metrics for monitoring and reporting purposes in the system.
* [**alerts.yaml**](/guides/connectors-developer-guide/required-files/alerts.yaml): This file contains the alert configuration, which will be converted into a Prometheus expression used to create the alert.


# metadata.json

This file is used by Obsrv while registering and running your connector. Here are the details of the file

| Name           | Description                                                                                                                                                                                                                                                                                                           |
| -------------- | --------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| `type`         | <p><strong>Type</strong>: <code>string</code><br><strong>Allowed Values</strong>: <code>"connector"</code><br>Specifies the type of the configuration. Must always be <code>"connector"</code>.</p>                                                                                                                   |
| `metadata`     | <p><strong>Type</strong>: <code>object</code><br>Contains detailed metadata for the connector, including attributes such as <code>id</code>, <code>version</code>, <code>tenant</code>, and others.</p>                                                                                                               |
| `id`           | <p><strong>Type</strong>: <code>string</code><br>A unique identifier for the connector.</p>                                                                                                                                                                                                                           |
| `name`         | <p><strong>Type</strong>: <code>string</code><br>The name of the connector.</p>                                                                                                                                                                                                                                       |
| `version`      | <p><strong>Type</strong>: <code>string</code><br>The version of the connector.</p>                                                                                                                                                                                                                                    |
| `tenant`       | <p><strong>Type</strong>: <code>string</code><br><strong>Allowed Values</strong>: <code>"single"</code>, <code>"multiple"</code><br>Indicates whether the connector is for a single or multiple tenants.</p>                                                                                                          |
| `type`         | <p><strong>Type</strong>: <code>string</code><br><strong>Allowed Values</strong>: <code>"source"</code><br>Specifies the type of connector, currently restricted to <code>"source"</code>.</p>                                                                                                                        |
| `category`     | <p><strong>Type</strong>: <code>string</code><br><strong>Allowed Values</strong>: <code>"batch"</code>, <code>"stream"</code><br>The processing category of the connector.</p>                                                                                                                                        |
| `description`  | <p><strong>Type</strong>: <code>string</code><br>A textual description of the connector.</p>                                                                                                                                                                                                                          |
| `technology`   | <p><strong>Type</strong>: <code>string</code><br><strong>Allowed Values</strong>: <code>"java"</code>, <code>"scala"</code>, <code>"python"</code><br>The primary technology used by the connector.</p>                                                                                                               |
| `runtime`      | <p><strong>Type</strong>: <code>string</code><br><strong>Allowed Values</strong>: <code>"spark"</code>, <code>"flink"</code><br>The runtime environment for the connector.<br><br>If category is <code>batch</code>, runtime should be <code>spark</code>; if <code>stream</code> it should be <code>flink</code></p> |
| `licence`      | <p><strong>Type</strong>: <code>string</code><br>Specifies the licence type of the connector.</p>                                                                                                                                                                                                                     |
| `owner`        | <p><strong>Type</strong>: <code>string</code><br>The owner or maintainer of the connector.</p>                                                                                                                                                                                                                        |
| `main_class`   | <p><strong>Type</strong>: <code>string</code><br>Specifies the main class of the program. Required for <code>java</code> or <code>scala</code> technologies, optional for <code>python</code>.</p>                                                                                                                    |
| `main_program` | <p><strong>Type</strong>: <code>string</code><br>Path to the main program of the connector.</p>                                                                                                                                                                                                                       |
| `icon`         | <p><strong>Type</strong>: <code>string</code> or <code>null</code><br>Path or URL for the icon representing the connector. Can be <code>null</code> if no icon is provided.</p>                                                                                                                                       |
| `connectors`   | <p><strong>Type</strong>: <code>array</code><br>A list of connectors associated with the configuration. Each item includes attributes like <code>id</code>, <code>name</code>, <code>description</code>, and <code>icon</code>.</p>                                                                                   |

Here is a sample `metadata.json` file for your reference

```json
{
  "type": "connector",
  "metadata": {
    "id": "example-connector",
    "name": "Example Connector",
    "description": "Pull data from a Source",
    "type": "source",
    "tenant": "single",
    "version": "1.0.0",
    "category": "stream",
    "technology": "scala",
    "runtime": "flink",
    "licence": "MIT",
    "owner": "Sunbird",
    "main_class": "org.sunbird.obsrv.connector.ExampleConnector",
    "main_program": "example-connector-1.0.0.jar",
    "icon": "icon.svg"
  }
}
```

Here is the JSON Schema to validate your `metadata.json`

```json
{
  "$schema": "https://json-schema.org/draft/2020-12/schema",
  "type": "object",
  "properties": {
    "type": {
      "type": "string",
      "enum": ["connector"]
    },
    "metadata": {
      "type": "object",
      "properties": {
        "id": {
          "type": "string"
        },
        "name": {
          "type": "string"
        },
        "version": {
          "type": "string"
        },
        "tenant": {
          "type": "string",
          "enum": ["single", "multiple"]
        },
        "type": {
          "type": "string",
          "enum": ["source"]
        },
        "category": {
          "type": "string",
          "enum": ["batch", "stream"]
        },
        "description": {
          "type": "string"
        },
        "technology": {
          "type": "string",
          "enum": ["java", "scala", "python"]
        },
        "runtime": {
          "type": "string",
          "enum": ["spark", "flink"]
        },
        "licence": {
          "type": "string"
        },
        "owner": {
          "type": "string"
        },
        "main_class": {
          "type": "string"
        },
        "main_program": {
          "type": "string"
        }
      },
      "required": [
        "id",
        "version",
        "tenant",
        "type",
        "category",
        "technology",
        "runtime",
        "licence",
        "owner",
        "main_program"
      ]
    },
    "connectors": {
      "type": "array",
      "items": {
        "type": "object",
        "properties": {
          "id": {
            "type": "string"
          },
          "name": {
            "type": "string"
          },
          "description": {
            "type": "string"
          },
          "icon": {
            "type": ["null", "string"]
          }
        },
        "required": [
          "id",
          "name",
          "description",
          "icon"
        ]
      }
    }
  },
  "required": [
    "type",
    "metadata"
  ],
  "dependentSchemas": {
    "metadata": {
      "allOf": [
        {
          "anyOf": [
            {
              "properties": {
                "metadata": {
                  "properties": {
                    "tenant": {
                      "enum": ["single"]
                    }
                  },
                  "required": [
                    "id",
                    "description",
                    "version",
                    "tenant",
                    "type",
                    "category",
                    "technology",
                    "runtime",
                    "licence",
                    "owner",
                    "icon",
                    "main_program"
                  ]
                }
              },
              "not": {
                "required": ["connectors"]
              }
            },
            {
              "properties": {
                "metadata": {
                  "properties": {
                    "tenant": {
                      "enum": ["multiple"]
                    }
                  },
                  "required": [
                    "id",
                    "name",
                    "version",
                    "tenant",
                    "type",
                    "category",
                    "technology",
                    "runtime",
                    "licence",
                    "owner",
                    "main_program"
                  ],
                  "not": {
                    "required": ["description"]
                  }
                },
                "connectors": {
                  "type": "array"
                }
              },
              "required": ["connectors"]
            }    
          ]
        },
        {
          "anyOf": [
            {
              "properties": {
                "metadata": {
                  "properties": {
                    "technology": {
                      "enum": ["python"]
                    }
                  },
                  "not": {
                     "required": ["main_class"]
                  }
                }
              }
            },
            {
              "properties": {
                "metadata": {
                  "properties": {
                    "technology": {
                      "enum": ["java", "scala"]
                    }
                  },
                  "required": ["main_class"]
                }
              }
            }
          ]
        },
        {
          "oneOf": [
            {
              "properties": {
                "metadata": {
                  "properties": {
                    "runtime": {
                      "enum": ["flink"]
                    },
                    "technology": {
                      "not": {
                        "enum": ["python"]
                      }
                    }
                  }
                }
              }
            },
            {
              "properties": {
                "metadata": {
                  "properties": {
                    "runtime": {
                      "enum": ["spark"]
                    },
                    "technology": {
                      "enum": ["python", "java", "scala"]
                    }
                  }
                }
              }
            }
          ]
        }
      ]
    }
  }
}
```


# ui-config.json

This page provides a comprehensive guide for developers on how to create and structure ui-config.json files used for auto-generating UIs for connectors.

### General Overview

1. **JSON Schema Standard**:
   * The file follows the **JSON Schema** standard (`type`, `properties`, `required`, etc.).
   * Developers should ensure compatibility with standard JSON Schema parsers for validation purposes.
2. **Custom Fields**:
   * `helptext`: A custom field providing detailed guidance that appears when a user focuses on a field.
   * `uiIndex`: Specifies the display order of fields in the UI.
3. **Flat Structure**: Ensure all fields are at the top level. Nested objects are not allowed.
4. **Key Attributes**:
   * `description`: Provides subtext for the input field.
   * `helptext`: Provides detailed guidance when the user focuses on the field.
   * `format`: Extended to include:
     * `hidden`: For fields with non-editable default values.
     * `password`: For sensitive inputs to hide entered text.
5. **Field Ordering**: Use the `uiIndex` attribute to dictate the order of fields in the UI.

***

### Core Sections in the `ui-config.json`

#### **1. Metadata**

* `title`: Title displayed in the UI.
* `description`: Brief summary of the connector setup.
* `helptext`: General instructions for the user.

```json
{
  "title": "Kafka Connector Setup Instructions",
  "description": "Configure Kafka Connector",
  "helptext": "Follow the below instructions to populate the required inputs needed for the connector correctly."
}
```

#### **2. Properties**

Each property represents a field in the UI. Define its type, validation rules, and additional UI hints.

**Example Field: Kafka Brokers**

```json
{
  "title": "Kafka Brokers",
  "type": "string",
  "description": "Enter Kafka broker addresses in the format: broker1-hostname:port,broker2-hostname:port",
  "helptext": "<p><strong>Kafka Broker Address Format:</strong> Enter the broker addresses in the following format:<code>&lt;broker1-hostname&gt;:&lt;port&gt;,&lt;broker2-hostname&gt;:&lt;port&gt;</code></p><p><em>Example:</em> <code>broker1.example.com:9092,broker2.example.com:9092</code></p>",
  "uiIndex": 1
}
```

#### **Field Attributes**

* **`type`**: Data type (e.g., `string`, `enum`).
* **`description`**: Short subtext visible below the field.
* **`helptext`**: Rich-text or HTML for focus-specific guidance.
* **`pattern`**: Regex for input validation.
* **`default`**: Default value, if any.
* **`enum`**: Allowed values for dropdowns or similar fields.
* **`format`**: Used for `hidden` or `password` to modify input behavior.

***

### Example Fields

Below are sample configurations for common Kafka connector fields:

1. **Kafka Brokers** (Required)
   * User inputs Kafka broker addresses.
   * Validation ensures correct format.
2. **Kafka Topic** (Required)
   * Alphanumeric with dots, dashes, or underscores.
   * Clear helptext and example.
3. **Kafka Auto Offset Reset** (Optional)
   * Dropdown with predefined values.
4. **Kafka Consumer Group ID** (Required)
   * Alphanumeric name identifying the consumer group.
5. **Data Format** (Hidden Default)
   * Predefined formats like `json`, `csv`, `parquet`.

***

### Required Fields

Specify the list of fields mandatory for the configuration:

```json
"required": [
  "source_kafka_broker_servers",
  "source_kafka_topic",
  "source_kafka_consumer_id",
  "source_kafka_auto_offset_reset",
  "source_data_format"
]
```

***

### Best Practices

1. **Descriptive Text**:
   * Ensure `description` and `helptext` are clear and concise.
   * Use examples where appropriate.
2. **UI Experience**:
   * Use `uiIndex` for logical field order.
   * Leverage `hidden` for technical defaults and `password` for secure inputs.
3. **Validation**:
   * Use `pattern` for format enforcement.
   * Provide feedback for invalid inputs.

***

### Example JSON

Here’s a full example for a Kafka connector:

```json
{
  "title": "Kafka Connector Setup Instructions",
  "description": "Configure Kafka Connector",
  "helptext": "Follow the instructions to configure the Kafka connector.",
  "type": "object",
  "properties": {
    "source_kafka_broker_servers": {
      "title": "Kafka Brokers",
      "type": "string",
      "description": "Enter Kafka broker addresses.",
      "helptext": "<p>Format: <code>broker1:port,broker2:port</code></p>",
      "uiIndex": 1
    },
    "source_kafka_topic": {
      "title": "Kafka Topic",
      "type": "string",
      "pattern": "^[a-zA-Z0-9\\\\._\\\\-]+$",
      "description": "Kafka topic name.",
      "helptext": "<p>Allowed characters: A-Z, 0-9, ., -, _</p>",
      "uiIndex": 2
    },
    "source_kafka_auto_offset_reset": {
      "title": "Kafka Auto Offset Reset",
      "type": "string",
      "enum": ["earliest", "latest", "none"],
      "default": "earliest",
      "description": "Auto offset reset.",
      "helptext": "Start reading from earliest, latest, or none.",
      "uiIndex": 3
    },
    "source_data_format": {
      "title": "Data Format",
      "type": "string",
      "enum": ["json", "csv"],
      "default": "json",
      "format": "hidden",
      "description": "Default format is JSON.",
      "uiIndex": 4
    }
  },
  "required": ["source_kafka_broker_servers", "source_kafka_topic"]
}
```

***

### Structure for Multi-Tenant Configurations

A multi-tenant connector's `ui-config.json` organizes configurations under source-specific keys. Each key contains the UI configuration specific to that source, using the same schema as single-source connectors.

**Example Structure**

Using object store connectors ui-config.json file:

```json
{
  "aws-s3": {
    "title": "AWS S3 Connector Setup",
    "description": "Configure AWS S3 for the connector",
    "helptext": "Enter details specific to AWS S3, such as bucket name and region.",
    "type": "object",
    "properties": {
      "s3_bucket": {
        "title": "S3 Bucket Name",
        "type": "string",
        "description": "Enter the name of the S3 bucket.",
        "helptext": "Specify the bucket name exactly as it appears in your AWS S3 console.",
        "uiIndex": 1
      },
      "s3_region": {
        "title": "S3 Region",
        "type": "string",
        "description": "Enter the AWS region of the S3 bucket.",
        "helptext": "Example: us-west-2",
        "uiIndex": 2
      },
      "s3_access_key": {
        "title": "Access Key",
        "type": "string",
        "description": "Enter your AWS Access Key ID.",
        "helptext": "This key is used to authenticate access to AWS S3.",
        "format": "password",
        "uiIndex": 3
      },
      "s3_secret_key": {
        "title": "Secret Key",
        "type": "string",
        "description": "Enter your AWS Secret Access Key.",
        "helptext": "This key is used to authenticate access to AWS S3.",
        "format": "password",
        "uiIndex": 4
      }
    },
    "required": ["s3_bucket", "s3_region", "s3_access_key", "s3_secret_key"]
  },
  "azure-blob": {
    "title": "Azure Blob Storage Setup",
    "description": "Configure Azure Blob Storage for the connector",
    "helptext": "Enter details specific to Azure Blob Storage, such as container name and connection string.",
    "type": "object",
    "properties": {
      "blob_container": {
        "title": "Container Name",
        "type": "string",
        "description": "Enter the Azure Blob container name.",
        "helptext": "Specify the container name exactly as it appears in your Azure Storage account.",
        "uiIndex": 1
      },
      "blob_connection_string": {
        "title": "Connection String",
        "type": "string",
        "description": "Enter your Azure Blob connection string.",
        "helptext": "Retrieve this from your Azure Portal under the storage account settings.",
        "format": "password",
        "uiIndex": 2
      }
    },
    "required": ["blob_container", "blob_connection_string"]
  }
}
```

***

### Guidelines for Multi-Tenant Connectors

1. **Source-Specific Keys**:
   * Use keys like to represent each source. For ex: an object store connector pulling data from multiple storage types can use the keys `aws-s3`, `azure-blob`, or `gcp-storage`
   * Each key encapsulates the configuration for that specific source type.
2. **Consistent Schema**:
   * The structure within each key follows the same JSON Schema standard as single-source connectors.
   * This ensures consistent validation and UI generation.
3. **UI Differentiation**:
   * Fields, help text, and validation can vary by source to accommodate source-specific requirements.
4. **Dynamic Selection**:
   * The UI can dynamically display the relevant configuration based on the source type selected by the user.

***

#### Best Practices for Multi-Tenant Configurations

1. **Clearly Defined Source Keys**:
   * Use meaningful and consistent source identifiers (e.g., `aws-s3`, `azure-blob`).
2. **Modular Design**:
   * Treat each source configuration as an independent module for easy updates and reusability.
3. **Dynamic UI**:
   * Leverage UI logic to dynamically display fields based on the selected `source_type`.
4. **Validation**:
   * Ensure each source configuration has clear validation rules and `required` fields.


# metrics.yaml

{% hint style="info" %}
Coming Soon!
{% endhint %}


# alerts.yaml

{% hint style="info" %}
Coming Soon!
{% endhint %}


# Obsrv Base Setup

To start developing connectors, here are the prerequisites

* Postgresql 16 or later
* Kafka 3.7.1 or later

## Postgresql

To set up a database for `obsrv` in PostgreSQL, follow these steps:

1. Start your PostgreSQL server.
2. Create a database named `obsrv` by executing the following command

```sql
CREATE DATABASE obsrv;
```

3. Proceed to create the required tables within the `obsrv` database as

```sql
CREATE TABLE public.datasets (
	id text NOT NULL,
	dataset_id text NULL,
	"type" text NOT NULL,
	"name" text NULL,
	validation_config json NULL,
	extraction_config json NULL,
	dedup_config json NULL,
	data_schema json NULL,
	denorm_config json NULL,
	router_config json NULL,
	dataset_config json NULL,
	tags _text NULL,
	data_version int4 NULL,
	status text NULL,
	created_by text NULL,
	updated_by text NULL,
	created_date timestamp DEFAULT now() NOT NULL,
	updated_date timestamp NOT NULL,
	published_date timestamp DEFAULT now() NOT NULL,
	api_version varchar(255) DEFAULT 'v1'::character varying NOT NULL,
	"version" int4 DEFAULT 1 NOT NULL,
	sample_data json DEFAULT '{}'::json NULL,
	entry_topic text DEFAULT 'dev.ingest'::text NOT NULL,
	CONSTRAINT datasets_pkey PRIMARY KEY (id)
);

CREATE TABLE public.connector_registry (
	id text NOT NULL,
	connector_id text NOT NULL,
	"name" text NOT NULL,
	"type" text NOT NULL,
	category text NOT NULL,
	"version" text NOT NULL,
	description text NULL,
	technology text NOT NULL,
	runtime text NOT NULL,
	licence text NOT NULL,
	"owner" text NOT NULL,
	iconurl text NULL,
	status text NOT NULL,
	ui_spec json DEFAULT '{}'::json NOT NULL,
	source_url text NOT NULL,
	"source" json NOT NULL,
	created_by text NOT NULL,
	updated_by text NOT NULL,
	created_date timestamp NOT NULL,
	updated_date timestamp NOT NULL,
	live_date timestamp NULL,
	CONSTRAINT connector_registry_connector_id_version_key UNIQUE (connector_id, version),
	CONSTRAINT connector_registry_pkey PRIMARY KEY (id)
);

CREATE TABLE public.connector_instances (
	id text NOT NULL,
	dataset_id text NOT NULL,
	connector_id text NOT NULL,
	connector_config text NOT NULL,
	operations_config json NOT NULL,
	status text NOT NULL,
	connector_state json DEFAULT '{}'::json NOT NULL,
	connector_stats json DEFAULT '{}'::json NOT NULL,
	created_by text NOT NULL,
	updated_by text NOT NULL,
	created_date timestamp NOT NULL,
	updated_date timestamp NOT NULL,
	published_date timestamp NOT NULL,
	"name" text NULL,
	CONSTRAINT connector_instances_pkey PRIMARY KEY (id),
	CONSTRAINT connector_instances_connector_id_fkey FOREIGN KEY (connector_id) REFERENCES public.connector_registry(id),
	CONSTRAINT connector_instances_dataset_id_fkey FOREIGN KEY (dataset_id) REFERENCES public.datasets(id)
);
```

{% hint style="info" %}
In case you already have a full version of Obsrv installed, the above steps can be skipped
{% endhint %}

## Kafka

The Obsrv pipeline utilizes Kafka to efficiently process data in real-time. Connector SDKs ensure data is pushed to Kafka upon confirmation of successful connection, facilitating smooth data integration and quick responsiveness.

Ensure that you have Kafka up and running


# Dev Requirements

{% tabs %}
{% tab title="Java / Scala" %}

## Prerequisites

### Stream Connectors

* Java 11
* Scala 2.12.11
* Apache Flink 1.17.2
* Apache Kafka 2.8.1

### Batch Connectors

* Java 11
* Scala 2.12.11
* Apache Spark 3.5.1
* Apache Kafka 2.8.1

## Libraries

Make sure you have the necessary repositories for the development

```sh
git clone git@github.com:Sunbird-Obsrv/job-sdk-scala.git
git clone git@github.com:Sunbird-Obsrv/connector-sdk-scala.git
```

### Setup

1. `job-sdk-scala`

```
cd job-sdk-scala
mvn clean install
```

2. `connector-sdk-scala`

```
cd connector-sdk-scala
mvn clean install
```

## Adding Dependencies

### Stream Connectors

Add the following to your project's `pom.xml` file under dependencies

{% code title="pom.xml" %}

```xml
<dependencies>
    ...
        <dependency>
            <groupId>org.sunbird.obsrv.connector</groupId>
            <artifactId>connector-sdk-flink</artifactId>
            <version>1.0.0</version>
        </dependency>
    ...
</dependencies>
```

{% endcode %}

### Batch Connectors

Add the following to your project's `pom.xml` file under dependencies

{% code title="pom.xml" %}

```xml
<dependencies>
    ...
        <dependency>
            <groupId>org.sunbird.obsrv.connector</groupId>
            <artifactId>connector-sdk-spark</artifactId>
            <version>1.0.0</version>
        </dependency>
    ...
</dependencies>
```

{% endcode %}
{% endtab %}

{% tab title="Python" %}

## Prerequisites

* Python 3.10 or higher
* Kafka 2.8.1
* Spark (PySpark) 3.5.1

## Required Packages

The `obsrv` python package is distributed through PyPI repository and can be installed using pip

```bash
pip install "obsrv[batch]"
```

### Using Poetry for Dependency Management

{% hint style="info" %}
`Poetry` is recommended for its ease of managing and packaging
{% endhint %}

[Poetry](https://python-poetry.org) is a popular tool for dependency management and packaging in Python projects. It streamlines the process of installing and updating project dependencies. To get started with Poetry, first install it using the following command

```bash
pip install poetry
```

Once installed, you can create a new Poetry project:

```bash
poetry new your_project_name
```

To add dependencies to your project, such as the `obsrv` package, use:

```bash
poetry add "obsrv[batch]"
```

Poetry automatically creates and manages a virtual environment for your project, ensuring isolated dependencies and compatibility management.
{% endtab %}
{% endtabs %}


# Interfaces

* [Stream Interfaces](/guides/connectors-developer-guide/interfaces/stream-interfaces)
* [Batch Interfaces](/guides/connectors-developer-guide/interfaces/batch-interfaces)


# Stream Interfaces

The SDK's expose an interface that has to be extended inorder to build a connector. A sample is shown below

{% tabs %}
{% tab title="Java" %}

## Imports

```java
package org.sunbird.obsrv.connector;

import com.typesafe.config.Config;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

import org.sunbird.obsrv.connector.model.Models;
import org.sunbird.obsrv.connector.source.IConnectorSource;
import org.sunbird.obsrv.connector.source.SourceConnectorFunction;
import org.sunbird.obsrv.job.exception.UnsupportedDataFormatException;

import java.util.List;
```

## SourceConnectorFunction

```java
public class ExampleSourceFunction extends SourceConnectorFunction {
    public ExampleSourceFunction(List<Models.ConnectorContext> connectorContexts) {
        super(connectorContexts);
    }

    @Override
    public void processEvent(
        String event, 
        Function1<String, BoxedUnit> onSuccess, 
        Function2<String, org.sunbird.obsrv.job.model.Models.ErrorData, BoxedUnit> onFailure, 
        Function2<String, Object, BoxedUnit> incMetric
    ){
        // TODO: Implement this method to process the event
        // Call onSuccess.apply(event) if the event is processed successfully
        // Call onFailure.apply(event, errorData) if the event processing fails
        // Call incMetric.apply(event, metricData) to increment the metric
    }

    @Override
    public List<String> getMetrics() {
        // TODO: Return the list of metrics
        return List.empty();
    }
}
```

Ref: [https://github.com/Sunbird-Obsrv/connector-sdk-scala/](https://github.com/Sunbird-Obsrv/connector-sdk-scala/blob/main/connector-sdk-flink/src/main/scala/org/sunbird/obsrv/connector/source/SourceConnectorFunction.scala#L10)

## IConnectorSource Class

```java
public class ExampleSourceConnector extends IConnectorSource {
    @Override
    public SingleOutputStreamOperator<String> getSourceStream(
        StreamExecutionEnvironment env, Config config
    ) throws UnsupportedDataFormatException {
        // TODO: Implement this method to return the source stream
        // env.fromSource(...)
    }

    @Override
    public SourceConnectorFunction getSourceFunction(
        List<Models.ConnectorContext> contexts, Config config) 
    {
        return ExampleSourceFunction(contexts)
    }
}
```

Ref: [https://github.com/Sunbird-Obsrv/connector-sdk-scala/](https://github.com/Sunbird-Obsrv/connector-sdk-scala/blob/main/connector-sdk-flink/src/main/scala/org/sunbird/obsrv/connector/source/IConnectorSource.scala)
{% endtab %}

{% tab title="Scala" %}

## Imports

```scala
package org.sunbird.obsrv.connector

import com.typesafe.config.Config
import org.apache.flink.api.common.serialization.SimpleStringSchema
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment
import org.json.{JSONException, JSONObject}
import org.sunbird.obsrv.connector.model.Models
import org.sunbird.obsrv.connector.source.{IConnectorSource, SourceConnector, SourceConnectorFunction}
import org.sunbird.obsrv.job.exception.UnsupportedDataFormatException
import org.sunbird.obsrv.job.model.Models.ErrorData
```

## SourceConnectorFunction

```scala

class ExampleSourceConnectorFunction(connectorContexts: List[ConnectorContext]) extends SourceConnectorFunction(connectorContexts) {

  /**
   * This method processes the incoming event.
   *
   * @param event The event to be processed.
   * @param onSuccess Callback function to be called on successful processing of the event.
   * @param onFailure Callback function to be called on failure in processing the event.
   * @param incMetric Function to increment the metric counter.
   */
  override def processEvent(event: String, 
                            onSuccess: String => Unit, 
                            onFailure: (String, ErrorData) => Unit, 
                            incMetric: (String, Long) => Unit): Unit = {
    // Implement your event processing logic here.
  }
  
  // TODO: Returns a list of custom metrics if any
  override def getMetrics(): List[String] = List[String]()
}
```

Ref: [https://github.com/Sunbird-Obsrv/connector-sdk-scala/](https://github.com/Sunbird-Obsrv/connector-sdk-scala/blob/main/connector-sdk-flink/src/main/scala/org/sunbird/obsrv/connector/source/SourceConnectorFunction.scala#L10)

## IConnectorSource

```scala
class ExampleConnectorSource extends IConnectorSource {

  @throws[UnsupportedDataFormatException]
  override def getSourceStream(env: StreamExecutionEnvironment, config: Config): SingleOutputStreamOperator[String] = {
    // Implement the logic to create and return the source stream
    // Example:
    // env.fromElements("event1", "event2", "event3")
  }

  override def getSourceFunction(contexts: List[ConnectorContext], config: Config): SourceConnectorFunction = {
    new ExampleSourceConnectorFunction(contexts)
  }
}
```

Ref: [https://github.com/Sunbird-Obsrv/connector-sdk-scala/](https://github.com/Sunbird-Obsrv/connector-sdk-scala/blob/main/connector-sdk-flink/src/main/scala/org/sunbird/obsrv/connector/source/IConnectorSource.scala)
{% endtab %}
{% endtabs %}


# Batch Interfaces

The SDK's exposes an interface which is to be extended inorder to build a connector

{% tabs %}
{% tab title="Java" %}

## Imports

```java
package org.sunbird.obsrv.connector;

import com.typesafe.config.Config;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;
import org.sunbird.obsrv.connector.model.Models.ConnectorContext;
import org.sunbird.obsrv.connector.source.ISourceConnector;

import java.util.Collections;
import java.util.Map;
```

## ISourceConnector

```java
public class ExampleSourceConnector implements ISourceConnector {

    @Override
    public Map<String, String> getSparkConf(Config config) {
        // TODO: Return the SparkConf related to your connector
        return Collections.emptyMap();
    }

    @Override
    public Dataset<Row> process(SparkSession spark, ConnectorContext ctx, Config config, BiConsumer<String, Long> metricFn) {
        // TODO: Add logic to read the data and return a dataframe
        return spark.emptyDataFrame();
    }
}
```

## Reference

* [https://github.com/Sunbird-Obsrv/connector-sdk-scala/](https://github.com/Sunbird-Obsrv/connector-sdk-scala/blob/main/connector-sdk-spark/src/main/scala/org/sunbird/obsrv/connector/source/ISourceConnector.scala)
  {% endtab %}

{% tab title="Scala" %}

## Imports

```scala
package org.sunbird.obsrv.connector

import com.typesafe.config.Config
import org.apache.spark.sql.functions.{col, max}
import org.apache.spark.sql.{DataFrame, Dataset, Row, SparkSession}
import org.sunbird.obsrv.connector.model.Models.ConnectorContext
import org.sunbird.obsrv.connector.source.{ISourceConnector, SourceConnector}
```

## ISourceConnector

```scala
class ExampleSourceConnector extends ISourceConnector {

  override def getSparkConf(config: Config): Map[String, String] = {
    // TODO: Return the SparkConf related to your connector
    Map[String, String]()
  }

  override def process(spark: SparkSession, ctx: ConnectorContext, config: Config, metricFn: (String, Long) => Unit): Dataset[Row] = {
    // TODO: Add logic to read the data and return a dataframe
    spark.emptyDataFrame
  }
```

## Reference

* [https://github.com/Sunbird-Obsrv/connector-sdk-scala/](https://github.com/Sunbird-Obsrv/connector-sdk-scala/blob/main/connector-sdk-spark/src/main/scala/org/sunbird/obsrv/connector/source/ISourceConnector.scala)
  {% endtab %}

{% tab title="Python" %}

## Imports

```python
from obsrv.common import ObsrvException
from obsrv.connector import ConnectorContext, MetricsCollector
from obsrv.connector.batch import ISourceConnector
from obsrv.job.batch import get_base_conf
from obsrv.models import ErrorData, StatusCode
from obsrv.utils import LoggerController
from pyspark.conf import SparkConf
from pyspark.sql import DataFrame, SparkSession
from pyspark.sql.functions import lit
```

## ISourceConnector

```python
class ExampleSource(ISourceConnector):
    def process(
        self,
        sc: SparkSession,
        ctx: ConnectorContext,
        connector_config: Dict[Any, Any],
        metrics_collector: MetricsCollector,
    ) -> Iterator[DataFrame]:
        # TODO: return or yield dataframe
        # yield sc.createDataFrame([], schema=None)
        return sc.createDataFrame([], schema=None)
        
    def get_spark_conf(self, connector_config) -> SparkConf:
        conf = get_base_conf()
        # TODO: Extend or Add the SparkConf related to your connector
        return conf

```

## Reference

* [https://github.com/Sunbird-Obsrv/obsrv-python-sdk](https://github.com/Sunbird-Obsrv/obsrv-python-sdk/blob/main/obsrv/connector/batch/source.py)
  {% endtab %}
  {% endtabs %}


# Classes


# ConnectorContext Class

## ConnectorContext

The [ConnectorContext](https://github.com/Sunbird-Obsrv/connector-sdk-scala/blob/main/connector-sdk-core/src/main/scala/org/sunbird/obsrv/connector/model/Models.scala) class is a part of the `org.sunbird.obsrv.connector.model` package. It encapsulates the context information for a connector instance, including its configuration, state, and statistics. This class is essential for managing the lifecycle and execution of connector instances within the Obsrv platform.

### Class Definition

```scala
package org.sunbird.obsrv.connector.model

import com.fasterxml.jackson.annotation.{JsonIgnore, JsonProperty}
import org.sunbird.obsrv.job.model.Models.ErrorData

object Models {

  case class ConnectorContext(
     @JsonProperty("connector_id") connectorId: String,
     @JsonProperty("dataset_id") datasetId: String,
     @JsonProperty("connector_instance_id") connectorInstanceId: String,
     @JsonProperty("connector_type") connectorType: String,
     @JsonIgnore entryTopic: String,
     @JsonIgnore state: ConnectorState,
     @JsonIgnore stats: ConnectorStats
   )
}
```

### Fields

#### `connectorId: String`

* **Description**: The unique identifier for the connector.
* **Annotations**: `@JsonProperty("connector_id")`

#### `datasetId: String`

* **Description**: The unique identifier for the dataset associated with the connector.
* **Annotations**: `@JsonProperty("dataset_id")`

#### `connectorInstanceId: String`

* **Description**: The unique identifier for the specific instance of the connector.
* **Annotations**: `@JsonProperty("connector_instance_id")`

#### `connectorType: String`

* **Description**: The type of the connector (e.g., source, sink).
* **Annotations**: `@JsonProperty("connector_type")`

#### `entryTopic: String`

* **Description**: The entry topic for the connector. This field is ignored during JSON serialization.
* **Annotations**: `@JsonIgnore`

#### [state: ConnectorState](https://github.com/Sunbird-Obsrv/connector-sdk-scala/blob/main/connector-sdk-core/src/main/scala/org/sunbird/obsrv/connector/model/ConnectorState.scala)

* **Description**: The state of the connector instance, encapsulated in a [ConnectorState](https://vscode-file/vscode-app/Applications/Visual%20Studio%20Code.app/Contents/Resources/app/out/vs/code/electron-sandbox/workbench/workbench.html) object. This field is ignored during JSON serialization.
* **Annotations**: `@JsonIgnore`

#### [stats: ConnectorStats](https://github.com/Sunbird-Obsrv/connector-sdk-scala/blob/main/connector-sdk-core/src/main/scala/org/sunbird/obsrv/connector/model/ConnectorStats.scala)

* **Description**: The statistics of the connector instance, encapsulated in a [ConnectorStats](https://vscode-file/vscode-app/Applications/Visual%20Studio%20Code.app/Contents/Resources/app/out/vs/code/electron-sandbox/workbench/workbench.html) object. This field is ignored during JSON serialization.
* **Annotations**: `@JsonIgnore`

### Usage

The ConnectorContext class is used to manage and track the state and statistics of a connector instance. It provides a structured way to access and manipulate the context information required for the execution of connectors.

#### Example

```scala
import org.sunbird.obsrv.connector.model.Models.ConnectorContext
import org.sunbird.obsrv.connector.model.{ConnectorState, ConnectorStats}
import org.sunbird.obsrv.job.util.PostgresConnectionConfig

implicit val postgresConfig: PostgresConnectionConfig = // initialize config

val state = new ConnectorState("connectorInstanceId", None)
val stats = new ConnectorStats("connectorInstanceId", None)

val context = ConnectorContext(
  connectorId = "connectorId",
  datasetId = "datasetId",
  connectorInstanceId = "connectorInstanceId",
  connectorType = "source",
  entryTopic = "entryTopic",
  state = state,
  stats = stats
)

// Accessing context fields
println(context.connectorId)
println(context.datasetId)
```

### See Also

* [ConnectorState](/guides/connectors-developer-guide/classes/connectorstate-class)
* [ConnectorStats](/guides/connectors-developer-guide/classes/connectorstats-class)


# ConnectorStats Class

## ConnectorStats

### Overview

The `ConnectorStats` class is a part of the `org.sunbird.obsrv.connector.model` package. It encapsulates the statistics information for a connector instance, allowing for the tracking and management of various metrics associated with the connector's performance.

### Class Definition

Ref: <https://github.com/Sunbird-Obsrv/connector-sdk-scala/blob/main/connector-sdk-core/src/main/scala/org/sunbird/obsrv/connector/model/ConnectorStats.scala>

### Arguments

#### `connectorInstanceId: String`

* **Description**: The unique identifier for the specific instance of the connector.

#### `statsJson: Option[String]`

* **Description**: The JSON representation of the statistics for the connector instance.

#### `stats: mutable.Map[String, AnyRef]`

* **Description**: A mutable map that holds the statistics metrics for the connector instance.

### Methods

#### `getStat[T](metric: String): Option[T]`

* **Description**: Retrieves the value of the specified metric.
* **Parameters**: `metric` - The name of the metric to retrieve.
* **Returns**: An `Option` containing the value of the metric, if it exists.

#### `getStat[T](metric: String, defaultValue: T): T`

* **Description**: Retrieves the value of the specified metric, or returns a default value if the metric does not exist.
* **Parameters**:
  * `metric` - The name of the metric to retrieve.
  * `defaultValue` - The default value to return if the metric does not exist.
* **Returns**: The value of the metric, or the default value.

#### `putStat[T <: AnyRef](metric: String, value: T): Unit`

* **Description**: Adds or updates the value of the specified metric.
* **Parameters**:
  * `metric` - The name of the metric to add or update.
  * `value` - The value to set for the metric.

#### `removeStat(metric: String): Option[AnyRef]`

* **Description**: Removes the specified metric from the statistics.
* **Parameters**: `metric` - The name of the metric to remove.
* **Returns**: An `Option` containing the removed value, if it existed.

#### `toJson(): String`

* **Description**: Serializes the statistics to a JSON string.
* **Returns**: A JSON string representation of the statistics.

#### `saveStats(): Unit`

* **Description**: Saves the current statistics to the database.
* **Throws**: `ObsrvException` if the statistics could not be saved.

### Usage

The `ConnectorStats` class is used to manage and track the performance metrics of a connector instance. It provides methods to access, update, and persist these metrics.

#### Example

```scala
import org.sunbird.obsrv.connector.model.ConnectorStats
import org.sunbird.obsrv.job.util.PostgresConnectionConfig

implicit val postgresConfig: PostgresConnectionConfig = // initialize config

val stats = new ConnectorStats("connectorInstanceId", None)

// Accessing and updating stats
stats.putStat("metric1", 100.asInstanceOf[AnyRef])
println(stats.getStat[Int]("metric1"))
stats.saveStats()
```


# ConnectorState Class

### Overview

The `ConnectorState` class is a part of the `org.sunbird.obsrv.connector.model` package. It encapsulates the state information for a connector instance, allowing for the tracking and management of various state attributes associated with the connector's lifecycle.

### Class Definition

Ref: <https://github.com/Sunbird-Obsrv/connector-sdk-scala/blob/main/connector-sdk-core/src/main/scala/org/sunbird/obsrv/connector/model/ConnectorState.scala>

### Args

#### `connectorInstanceId: String`

* **Description**: The unique identifier for the specific instance of the connector.

#### `stateJson: Option[String]`

* **Description**: The JSON representation of the state for the connector instance.

#### `state: mutable.Map[String, AnyRef]`

* **Description**: A mutable map that holds the state attributes for the connector instance.

### Methods

#### `getState[T](attribute: String): Option[T]`

* **Description**: Retrieves the value of the specified state attribute.
* **Parameters**: `attribute` - The name of the state attribute to retrieve.
* **Returns**: An `Option` containing the value of the state attribute, if it exists.

#### `getState[T](attribute: String, defaultValue: T): T`

* **Description**: Retrieves the value of the specified state attribute, or returns a default value if the attribute does not exist.
* **Parameters**:
  * `attribute` - The name of the state attribute to retrieve.
  * `defaultValue` - The default value to return if the attribute does not exist.
* **Returns**: The value of the state attribute, or the default value.

#### `putState[T <: AnyRef](attrib: String, value: T): Unit`

* **Description**: Adds or updates the value of the specified state attribute.
* **Parameters**:
  * `attrib` - The name of the state attribute to add or update.
  * `value` - The value to set for the state attribute.

#### `removeState(attrib: String): Option[AnyRef]`

* **Description**: Removes the specified state attribute from the state.
* **Parameters**: `attrib` - The name of the state attribute to remove.
* **Returns**: An `Option` containing the removed value, if it existed.

#### `contains(attrib: String): Boolean`

* **Description**: Checks if the specified state attribute exists in the state.
* **Parameters**: `attrib` - The name of the state attribute to check.
* **Returns**: `true` if the attribute exists, `false` otherwise.

#### `toJson(): String`

* **Description**: Serializes the state to a JSON string.
* **Returns**: A JSON string representation of the state.

#### `saveState(): Unit`

* **Description**: Saves the current state to the database.
* **Throws**: `ObsrvException` if the state could not be saved.

### Usage

The `ConnectorState` class is used to manage and track the state attributes of a connector instance. It provides methods to access, update, and persist these attributes.

#### Example

```scala
import org.sunbird.obsrv.connector.model.ConnectorState
import org.sunbird.obsrv.job.util.PostgresConnectionConfig

implicit val postgresConfig: PostgresConnectionConfig = // initialize config

val state = new ConnectorState("connectorInstanceId", None)

// Accessing and updating state
state.putState("attribute1", "value1".asInstanceOf[AnyRef])
println(state.getState[String]("attribute1"))
state.saveState()
```


# ErrorData Class

### Overview

The `ErrorData` case class is a part of the `org.sunbird.obsrv.job.model.Models` package. It encapsulates error information, including an error code and an error message, which can be used to provide detailed error reporting and handling within the application.

### Class Definition

```scala
package org.sunbird.obsrv.job.model.Models

import com.fasterxml.jackson.annotation.JsonProperty

case class ErrorData(
  @JsonProperty("error_code") errorCode: String,
  @JsonProperty("error_msg") errorMsg: String
)
```

### Fields

#### `errorCode: String`

* **Description**: The unique code representing the specific error.
* **Annotations**: `@JsonProperty("error_code")`

#### `errorMsg: String`

* **Description**: A descriptive message providing details about the error.
* **Annotations**: `@JsonProperty("error_msg")`

### Usage

The `ErrorData` case class is used to encapsulate error information that can be serialized and deserialized to and from JSON. It provides a structured way to represent errors in the application.

#### Example

```scala
import org.sunbird.obsrv.job.model.Models.ErrorData

val error = ErrorData(
  errorCode = "404",
  errorMsg = "Resource not found"
)

// Accessing error fields
println(s"Error Code: ${error.errorCode}")
println(s"Error Message: ${error.errorMsg}")
```

### JSON Representation

An instance of `ErrorData` can be serialized to JSON as follows:

```json
{
  "error_code": "404",
  "error_msg": "Resource not found"
}
```


# MetricData Class

### Overview

The `MetricData` case class is a part of the `org.sunbird.obsrv.job.model.Models` package. It encapsulates metric information, including a map of metric names to their values and a list of labels associated with these metrics. This class is useful for representing and managing metrics data within the application.

### Class Definition

```scala
package org.sunbird.obsrv.job.model.Models

case class MetricData(
  metric: Map[String, Long],
  labels: List[Map[String, String]]
)
```

### Fields

#### `metric: Map[String, Long]`

* **Description**: A map where the keys are metric names and the values are the corresponding metric values.
* **Type**: `Map[String, Long]`

#### `labels: List[Map[String, String]]`

* **Description**: A list of maps, where each map represents a set of labels associated with the metrics.
* **Type**: `List[Map[String, String]]`

### Usage

The `MetricData` case class is used to encapsulate and manage metrics data, including the metric values and their associated labels. It provides a structured way to represent metrics in the application.

#### Example

```scala
import org.sunbird.obsrv.job.model.Models.MetricData

val metric = Map("metric1" -> 100L, "metric2" -> 200L)
val labels = List(
  Map("label1" -> "value1", "label2" -> "value2"),
  Map("label3" -> "value3", "label4" -> "value4")
)

val metricData = MetricData(metric, labels)

// Accessing metric data fields
println(s"Metrics: ${metricData.metric}")
println(s"Labels: ${metricData.labels}")
```

### JSON Representation

An instance of `MetricData` can be serialized to JSON as follows:

```json
{
  "metric": {
    "metric1": 100,
    "metric2": 200
  },
  "labels": [
    {
      "label1": "value1",
      "label2": "value2"
    },
    {
      "label3": "value3",
      "label4": "value4"
    }
  ]
}
```


# Verifying

## Environment Setup

1. Since connectors are configured at a dataset level, create a sample dataset. Here is a sample to create `new-york-taxi-data` dataset.

```sql
INSERT INTO public.datasets (id,dataset_id,"type","name",validation_config,extraction_config,dedup_config,data_schema,denorm_config,router_config,dataset_config,tags,data_version,status,created_by,updated_by,created_date,updated_date,published_date,api_version,"version",sample_data,entry_topic) VALUES
 ('new-york-taxi-data','new-york-taxi-data','event','new-york-taxi-data','{"validate": true, "mode": "Strict"}','{"is_batch_event": true, "batch_id": "id", "extraction_key": "events", "dedup_config": {"dedup_key": "id", "drop_duplicates": true, "dedup_period": 604800}}','{"dedup_key": "tripID", "drop_duplicates": true, "dedup_period": 604800}','{"$schema": "https://json-schema.org/draft/2020-12/schema", "type": "object", "properties": {"tripID": {"type": "string", "suggestions": [{"message": "The Property ''''tripID'''' appears to be ''''uuid'''' format type.", "advice": "Suggest to not to index the high cardinal columns", "resolutionType": "DEDUP", "severity": "LOW", "path": "properties.tripID"}], "arrival_format": "text", "data_type": "string"}, "VendorID": {"type": "string", "arrival_format": "text", "data_type": "string"}, "tpep_pickup_datetime": {"type": "string", "suggestions": [{"message": "The Property ''''tpep_pickup_datetime'''' appears to be ''''date-time'''' format type.", "advice": "The System can index all data on this column", "resolutionType": "INDEX", "severity": "LOW", "path": "properties.tpep_pickup_datetime"}], "arrival_format": "text", "data_type": "date-time"}, "tpep_dropoff_datetime": {"type": "string", "suggestions": [{"message": "The Property ''''tpep_dropoff_datetime'''' appears to be ''''date-time'''' format type.", "advice": "The System can index all data on this column", "resolutionType": "INDEX", "severity": "LOW", "path": "properties.tpep_dropoff_datetime"}], "arrival_format": "text", "data_type": "date-time"}, "passenger_count": {"type": "string", "arrival_format": "text", "data_type": "string"}, "trip_distance": {"type": "string", "arrival_format": "text", "data_type": "string"}, "RatecodeID": {"type": "string", "arrival_format": "text", "data_type": "string"}, "store_and_fwd_flag": {"type": "string", "arrival_format": "text", "data_type": "string"}, "PULocationID": {"type": "string", "arrival_format": "text", "data_type": "string"}, "DOLocationID": {"type": "string", "arrival_format": "text", "data_type": "string"}, "payment_type": {"type": "string", "arrival_format": "text", "data_type": "string"}, "primary_passenger": {"type": "object", "properties": {"email": {"type": "string", "arrival_format": "text", "data_type": "string"}, "mobile": {"type": "string", "arrival_format": "text", "data_type": "string"}}, "arrival_format": "object", "data_type": "object", "additionalProperties": false}, "fare_details": {"type": "object", "properties": {"fare_amount": {"type": "string", "arrival_format": "text", "data_type": "string"}, "extra": {"type": "string", "arrival_format": "text", "data_type": "string"}, "mta_tax": {"type": "string", "arrival_format": "text", "data_type": "string"}, "tip_amount": {"type": "string", "arrival_format": "text", "data_type": "string"}, "tolls_amount": {"type": "string", "arrival_format": "text", "data_type": "string"}, "improvement_surcharge": {"type": "string", "arrival_format": "text", "data_type": "string"}, "total_amount": {"type": "string", "arrival_format": "text", "data_type": "string"}, "congestion_surcharge": {"type": "string", "arrival_format": "text", "data_type": "string"}}, "arrival_format": "object", "data_type": "object", "additionalProperties": false}}, "additionalProperties": false}','{"denorm_fields": [], "redis_db_host": "localhost", "redis_db_port": 6379}','{"topic": "new-york-taxi-data"}','{"keys_config": {"timestamp_key": "obsrv_meta.syncts", "data_key": "", "partition_key": ""}, "indexing_config": {"olap_store_enabled": true, "lakehouse_enabled": false, "cache_enabled": false}, "cache_config": {"redis_db_host": "localhost", "redis_db_port": 6379, "redis_db": 0}, "file_upload_path": []}','{}',2,'Live','SYSTEM','SYSTEM','2024-09-12 13:13:15.400866','2024-09-17 11:26:00.774249','2024-09-17 11:26:00.774249','v2',1,'{"tripID": "0de066f6-e7e4-44e3-9d5a-af8be8cd360a", "VendorID": "2", "tpep_pickup_datetime": "2023-04-12 13:48:30", "tpep_dropoff_datetime": "2023-11-17 13:52:40", "passenger_count": "3", "trip_distance": ".00", "RatecodeID": "1", "store_and_fwd_flag": "N", "PULocationID": "236", "DOLocationID": "236", "payment_type": "1", "primary_passenger": {"email": "Felipe.Grant@hotmail.com", "mobile": "494-699-7052"}, "fare_details": {"fare_amount": "4.5", "extra": "0.5", "mta_tax": "0.5", "tip_amount": "0", "tolls_amount": "0", "improvement_surcharge": "0.3", "total_amount": "5.8", "congestion_surcharge": ""}, "suggestedPii": []}','local.ingest');
```

2. In order to run a connector, the connector must me registered in Obsrv. To do so, fill in the values and run the below query

```sql
INSERT INTO public.connector_registry (id,connector_id,"name","type",category,"version",description,technology,runtime,licence,"owner",iconurl,status,ui_spec,source_url,"source",created_by,updated_by,created_date,updated_date,live_date) VALUES
	 ('example-connector',
	 'example-connector',
	 'Example Connector',
	 'source',
	 'stream',
	 '1.0.0',
	 'Pull data from a Source',
	 'scala',
	 'flink',
	 'MIT',
	 'Sunbird',
	 'data:image/svg+xml;base64,,',
	 'Live','{}','connector-1.0.0-distribution.tar.gz','{}',
	 'SYSTEM','SYSTEM',
	 now(), now(), now()
	);
```

3. The connector must have an active instance inorder where we specify the mapping of the dataset and the connector created above.

```sql
INSERT INTO public.connector_instances (id,dataset_id,connector_id,connector_config,operations_config,status,connector_state,connector_stats,created_by,updated_by,created_date,updated_date,published_date) VALUES
	 ('example-instance-1','new-york-taxi-data','example-connector-1.0.0',
	 '<aes-256-ecb encrypted string>',
	 '{}','Live','{}','{}','SYSTEM','SYSTEM', now(), now(), now()
	 );
```

{% hint style="info" %}
The `connector_config` in the instances uses AES ECB encryption with a key size of 256 bits.
{% endhint %}

> To verify if the encryption is correct use an online tool like [this](https://encode-decode.com/aes-256-ecb-encrypt-online/) where algorithm is **AES**, mode is **ECB**, **No-Padding** and key size is **256 bits** (32 character).
>
> For example the following JSON Config\
> `{"source_ip": "localhost", "source_port": 5432, "source_db": "obsrv", "table": "datasets"}`<br>
>
> with a secret `strong_encryption_key_to_encrypt`<br>
>
> should generate an encrypted string equal to
>
> `mmT3tqwpq798ywgsdfEdhtp5VXcFzFjutmaEuFqAw9fl4G4OtK6vz89Nhd/deChFWzW3yJWB3Y//SpDzwc3RY31FGqwnAhhD3fpjXP+XYHthCWxBIKar5j8pdTBZ867J`<br>
>
> which can be verfied online, before you encrypt your config. Upon decrypting it should be a valid JSON

4. Make sure you have the configuration files based on the language you are writing the connector in

{% tabs %}
{% tab title="Java / Scala" %}
Here is a sample file to be used in Java / Scala

{% code title="obsrv-connector-config.conf" %}

```editorconfig
postgres {
  host = localhost
  port = 5432
  maxConnections = 2
  user = "postgres"
  password = "postgres"
  database = "obsrv"
}

kafka {
  producer = {
    broker-servers = "localhost:9092"
    compression = "snappy"
    max-request-size = 1000000 # 1MB
  }
  output = {
    connector = {
      failed.topic = "connector.failed"
      metric.topic = "obsrv-connectors-metrics"
    }
  }
}

obsrv.encryption.key = "strong_encryption_key_to_encrypt"

task {
  checkpointing.compressed = true
  checkpointing.interval = 60000
  checkpointing.pause.between.seconds = 30000
  restart-strategy.attempts = 3
  restart-strategy.delay = 30000 # in milli-seconds
  parallelism = 1
  consumer.parallelism = 1
  downstream.operators.parallelism = 1
}

env = local
building-block = obsrv
```

{% endcode %}

Please make sure you update the database and kafka connection values as per your setup
{% endtab %}

{% tab title="Python" %}
Here is a sample configuration file to be used in python

{% code title="obsrv-connector-config.yaml" %}

```yaml
postgres:
  dbname: obsrv
  user: postgres
  password: postgres
  host: localhost
  port: 5432

kafka:
  broker-servers: localhost:9092
  telemetry-topic: obsrv-connectors-telemetry
  connector-metrics-topic: obsrv-connectors-metrics
  producer:
    compression: snappy
    max-request-size: 1000000

obsrv_encryption_key: strong_encryption_key_to_encrypt

connector_instance_id: example-instance-1

building-block: obsrv
env: local
```

{% endcode %}

Please make sure you update the database and kafka connection values as per your setup
{% endtab %}
{% endtabs %}

## Running Stream Connectors

1. Make sure you have Apache Flink running and is ready to accept jobs for processing. Make sure you have atleast ONE task slot for submitting the job
2. Submit the JAR using flink submit.

{% code overflow="wrap" %}

```bash
$FLINK_HOME/bin/flink run <path_to_jar> --port <random_port> --config.file.path <path_to_obsrv-connector-config.conf> --metadata.id <id_from_connector_registry>
```

{% endcode %}

## Running Batch Connectors

1. Make sure you have spark installed
2. Submit the JAR to Spark using the following command

{% tabs %}
{% tab title="Java / Scala" %}
{% code overflow="wrap" %}

```bash
spark-submit --master=local[*] --jars <include_dependent_jars> --class <main_class> <path_to_jar> -f <path_to_obsrv-connector-config.conf> -c <id_from_connector_instances>
```

{% endcode %}
{% endtab %}

{% tab title="Python" %}
{% code overflow="wrap" %}

```bash
spark-submit --master=local\[\*\] --conf "spark.pyspark.driver.python=<path_to_python_executable>" --conf "spark.pyspark.python=<path_to_python_executable>" --jars <include_dependent_jars> <path_to_main_python_file> -f <path_to_obsrv-connector-config.yaml> -c <id_from_connector_instances>
```

{% endcode %}

{% hint style="info" %}
If using an virtual environment for python development, then `path_to_python_executable` must be from the same environment
{% endhint %}
{% endtab %}
{% endtabs %}


# Packaging Guide

## Bundling Streaming Connectors

### Java / Scala

{% hint style="info" %}
The below document outlines the structure if you are using `mvn` as your build tool. If you are using `sbt` or others, use the following as reference and update accordingly to generate the package structure.
{% endhint %}

The stream connectors are expected to be bundled as a Single JAR (Fat JAR). We recommend using some thing like `maven-shade-plugin` to build the final JAR file.

Additionally, we recommend using `maven-assembly-plugin` to bundle all required files and JAR into a single distribution.

Here is a sample using `maven-shade-plugin` and `maven-assembly-plugin`

{% code title="pom.xml" %}

```xml
<build>
    ...
    <plugins>
        ...
        <plugin>
            <groupId>org.apache.maven.plugins</groupId>
            <artifactId>maven-shade-plugin</artifactId>
            <version>3.2.1</version>
            <executions>
                <execution>
                    <phase>package</phase>
                    <goals>
                        <goal>shade</goal>
                    </goals>
                    <configuration>
                        <shadedArtifactAttached>false</shadedArtifactAttached>
                        <artifactSet>
                            <excludes>
                                <exclude>com.google.code.findbugs:jsr305</exclude>
                            </excludes>
                        </artifactSet>
                        <filters>
                            <filter>
                                <!-- Do not copy the signatures in the META-INF folder.
                                Otherwise, this might cause SecurityExceptions when using the JAR. -->
                                <artifact>*:*</artifact>
                                <excludes>
                                    <exclude>META-INF/*.SF</exclude>
                                    <exclude>META-INF/*.DSA</exclude>
                                    <exclude>META-INF/*.RSA</exclude>
                                </excludes>
                            </filter>
                        </filters>
                        <transformers>
                            <transformer
                                    implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer">
                                <mainClass>org.sunbird.obsrv.connector.ExampleSourceConnector</mainClass>
                            </transformer>
                            <!-- append default configs -->
                            <transformer
                                    implementation="org.apache.maven.plugins.shade.resource.AppendingTransformer">
                                <resource>reference.conf</resource>
                            </transformer>
                        </transformers>
                    </configuration>
                </execution>
            </executions>
        </plugin>
        <plugin>
            <artifactId>maven-assembly-plugin</artifactId>
            <executions>
                <execution>
                    <id>distro-assembly</id>
                    <phase>package</phase>
                    <goals>
                        <goal>single</goal>
                    </goals>
                    <configuration>
                        <descriptors>
                            <descriptor>src/main/assembly/src.xml</descriptor>
                        </descriptors>
                    </configuration>
                </execution>
            </executions>
        </plugin>
        ...
    </plugins>
</build>
```

{% endcode %}

And here is a sample of `src/main/assembly/src.xml` for building the distribution

{% code title="src.xml" %}

```xml
<?xml version="1.0" encoding="UTF-8"?>

<assembly xmlns="http://maven.apache.org/ASSEMBLY/2.2.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
          xsi:schemaLocation="http://maven.apache.org/ASSEMBLY/2.2.0 https://maven.apache.org/xsd/assembly-2.2.0.xsd">
    <id>distribution</id>
    <formats>
        <format>tar.gz</format>
    </formats>
    <fileSets>
        <fileSet>
            <directory>${basedir}/src/main/resources</directory>
            <outputDirectory>/</outputDirectory>
            <includes>
                <include>*</include>
            </includes>
            <excludes>
                <exclude>obsrv-connector-config.conf</exclude>
            </excludes>
        </fileSet>
        <fileSet>
            <directory>target</directory>
            <outputDirectory>/</outputDirectory>
            <includes>
                <include>*.jar</include>
            </includes>
        </fileSet>
    </fileSets>
</assembly>
```

{% endcode %}

Run `mvn clean package` to generate the distribution

The contents of the distribution should be in the following format

```
example-connector-0.1.0-distribution.tar.gz
├── example_connector.jar
├── alerts.yaml
├── metadata.json
├── metrics.yaml
├── ui-config.json
└── icon.svg
```

## Bundling Batch Connectors

### Java / Scala

The batch connectors are bundled as a standalone JAR file, where the dependent JARs are to be included in the `libs` folder in the bundle, by using the `maven-assembly-plugin.`

Here is a sample of including the `maven-assembly-plugin` to your build

{% code title="pom.xml" %}

```xml
<build>
    ...
    <plugins>
        ...
        <plugin>
            <artifactId>maven-assembly-plugin</artifactId>
            <executions>
                <execution>
                    <id>distro-assembly</id>
                    <phase>package</phase>
                    <goals>
                        <goal>single</goal>
                    </goals>
                    <configuration>
                        <descriptors>
                            <descriptor>src/main/assembly/src.xml</descriptor>
                        </descriptors>
                    </configuration>
                </execution>
            </executions>
        </plugin>
        ...
    </plugins>
</build>
```

{% endcode %}

And here is a sample of `src/main/assembly/src.xml` for building the distribution

{% code title="src.xml" %}

```xml
<?xml version="1.0" encoding="UTF-8"?>

<assembly xmlns="http://maven.apache.org/ASSEMBLY/2.2.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
          xsi:schemaLocation="http://maven.apache.org/ASSEMBLY/2.2.0 https://maven.apache.org/xsd/assembly-2.2.0.xsd">
    <id>distribution</id>
    <formats>
        <format>tar.gz</format>
    </formats>
    <fileSets>
        <fileSet>
            <directory>${basedir}/src/main/resources</directory>
            <outputDirectory>/</outputDirectory>
            <includes>
                <include>*</include>
            </includes>
            <excludes>
                <exclude>obsrv-connector-config.conf</exclude>
            </excludes>
        </fileSet>
        <fileSet>
            <directory>target</directory>
            <outputDirectory>/</outputDirectory>
            <includes>
                <include>*.jar</include>
            </includes>
        </fileSet>
    </fileSets>
    <dependencySets>
        <dependencySet>
            <outputDirectory>libs</outputDirectory>
            <useTransitiveFiltering>true</useTransitiveFiltering>
        </dependencySet>
    </dependencySets>
</assembly>
```

{% endcode %}

Run `mvn clean package` to generate the distribution

The contents of the distribution should be in the following format

```
example-connector-0.1.0-distribution.tar.gz
├── libs
    └── sample-dependency.jar
├── example_connector.jar
├── alerts.yaml
├── metadata.json
├── metrics.yaml
├── requirements.txt
├── ui-config.json
└── icon.svg
```

### Python

Since we use PySpark for building the batch connectors in Python, we are required to include the dependent JARs as a part of the distribution.

Here are the steps if you are using `poetry` as a package manager for Python Connectors.

Add a `build_dist.py` script to the `scripts` folder with the following contents. Briefly the script does the following

* It exports the python requirements to `requirements.txt` file, so that these can be installed on the runtime.
* It downloads all the dependent JARs that are to be included in the package using `mvn`

{% code title="scripts/build\_dist.py" %}

```python
import subprocess
import os

def main():
    # remove dirs recursively
    subprocess.Popen("rm -rf dist", shell=True).wait()

    # Path to the directory containing the JAR files
    jar_dir = os.path.join(os.path.dirname(__file__), '..', 'libs')
    os.makedirs(jar_dir, exist_ok=True)

    subprocess.Popen("""poetry export --without-hashes --format=requirements.txt | awk '{split($0,a,"; "); print a[1]}' > requirements.txt""", shell=True).wait()
    # TODO: Only uncomment the below line, if you have Java dependencies for your script, which have to be included in the libs folder
    # subprocess.Popen("mvn dependency:copy-dependencies -DrepoUrl=http://repo1.maven.org/maven2/ -DexcludeTrans -DoutputDirectory=libs", shell=True).wait()
    subprocess.Popen("poetry build -f sdist", shell=True).wait()
    subprocess.Popen("rm -rf requirements.txt libs", shell=True).wait()

if __name__ == '__main__':
    main()
```

{% endcode %}

The packaging instructions have to be included in the `pyproject.toml` file.

```toml
packages = [
    { include = "example_connector", from = ".", format = "sdist" },
]

include = [
    "requirements.txt",
    # "libs/*.jar", # TODO: Only uncomment the below line, if you have Java dependencies for your script, which have to be included in the libs folder
    "ui-config.json",
    "metadata.json",
    "alerts.yaml",
    "metrics.yaml",
    "icon.svg" # Add additional files like the icons that are to be included in the bundle
]

[build-system]
requires = ["poetry-core"]
build-backend = "poetry.core.masonry.api"

[tool.poetry.scripts]
package = "scripts.build_dist:main"
```

Add a `pom.xml` file to the root and specify the dependent JARs required by the connector

{% code title="pom.xml" %}

```xml
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
    xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>
    <groupId>in.sanketika.connectors</groupId>
    <artifactId>example-connector</artifactId>
    <version>1.0.0</version>
    <dependencies>
        <dependency>
            <groupId> </groupId>
            <artifactId> </artifactId>
            <version> </version>
        </dependency>
        ...
    </dependencies>
</project>
```

{% endcode %}

Run `poetry run package` to build the distribution that can be installed on Obsrv.

The contents of the distribution should be in the following format

```
example-connector-0.1.0-distribution.tar.gz
├── libs
    └── sample-dependency.jar
├── example_connector
    └── main.py
├── alerts.yaml
├── metadata.json
├── metrics.yaml
├── requirements.txt
├── ui-config.json
└── icon.svg
```


# Reference Implementations

## Stream Connectors

### Scala

* [Kafka Connector](https://github.com/Sunbird-Obsrv/kafka-connector)

## Batch Connectors

### Scala

* [JDBC Connector](https://github.com/Sunbird-Obsrv/jdbc-connector)

### Python

* [Object Store Connector](https://github.com/Sunbird-Obsrv/object-store-connector)


# Coming Soon!

List of guides that will be added soon:

* How to connect superset with Obsrv and create dashboards in superset
* How to configure and access monitoring dashboards
* How to scale up and scale down Obsrv and its components
* How to configure alerts & notifications
* How to uninstall/remove optional components in an Obsrv installation


# Community

How to participate in Obsrv community

### Welcome to our Community Section! <a href="#id-2ltvlps1n48x" id="id-2ltvlps1n48x"></a>

Welcome to our community section dedicated to leveraging the power of Obsrv core capabilities, services and technologies used. This document serves as a guide for community members on how to effectively use our GitHub, Discord & Jira platforms to contribute to our projects and engage with the community.

### Getting started <a href="#id-5tj2bu9nts5n" id="id-5tj2bu9nts5n"></a>

### 1. Explore Our GitHub Projects and Discussions <a href="#rs02d61kgjt" id="rs02d61kgjt"></a>

If you haven't already, Explore our GitHub [repositories](https://github.com/orgs/Sunbird-Obsrv/repositories?q=obsrv\&type=all\&language=\&sort=) and [Discussions](https://github.com/orgs/Sunbird-Obsrv/discussions). This will give you access to our repositories and allow you to contribute to our projects and discussions.

### 2. Explore Our Projects on Jira <a href="#id-6h2jvs27rfmy" id="id-6h2jvs27rfmy"></a>

Visit our Jira [workspace](https://project-sunbird.atlassian.net/jira/software/c/projects/OB/issues) to explore our projects, track issues, and contribute to discussions.

### 3. Join our Group on Discord <a href="#ovab1v6x3rfo" id="ovab1v6x3rfo"></a>

Join our discord [channel](https://discord.gg/Q5mvw2mGC8) to initiate or contribute to any discussions and raise/report issues.

### How to Contribute <a href="#sdm7g8g05qu1" id="sdm7g8g05qu1"></a>

### 1. Forking and Cloning Repositories <a href="#aad92vaqbbyg" id="aad92vaqbbyg"></a>

Fork the repository you want to contribute to on GitHub, then clone it to your local machine using Git. This allows you to work on the codebase locally.

### 2. Making Changes <a href="#i1ybrw42w9du" id="i1ybrw42w9du"></a>

Make the necessary changes or additions to the codebase on your local machine. Ensure that your changes adhere to our coding standards and guidelines.

### 3. Submitting Pull Requests <a href="#udjumoxoo1er" id="udjumoxoo1er"></a>

Once you're done making changes, push your commits to your forked repository on GitHub and submit a pull request to the original repository. Provide a clear description of your changes and reference any related issues.

### 4. Participating in Discussions <a href="#id-7s5nmcanzlh7" id="id-7s5nmcanzlh7"></a>

Engage with the community by participating in discussions on GitHub discussions, Discord channel or Jira tickets. Share your ideas, ask questions, and provide feedback.

### Connect with Us <a href="#id-1jcdy691svj" id="id-1jcdy691svj"></a>

Stay updated with the latest news and announcements by following our [announcements](https://github.com/orgs/Sunbird-Obsrv/discussions/categories/announcements).

### Conclusion <a href="#op06p7rh5rj4" id="op06p7rh5rj4"></a>

Thank you for being a part of our community and contributing to our projects. Together, we can achieve great things through collaboration, innovation, and teamwork!


# Previous Versions


# SB-5.0 Version

Older version of Obsrv released as part of Sunbird 5.0. This version will be deprecated soon.


# Overview

Older version of Obsrv released as part of Sunbird 5.0. This version will be deprecated soon.


# USE


# Release Notes

In the subsequent pages you will find detailed documentation of these releases:

### Obsrv 2.0 Release Notes

{% content-ref url="<https://github.com/Sunbird-Obsrv/Community/tree/main/previous-versions/sb-5.0-version/use/release-notes/obsrv-2.0.1-ga.md>" %}
<https://github.com/Sunbird-Obsrv/Community/tree/main/previous-versions/sb-5.0-version/use/release-notes/obsrv-2.0.1-ga.md>
{% endcontent-ref %}

{% content-ref url="/pages/W18TT4twEO4odATmO9XN" %}
[Obsrv 2.0.0-GA](/previous-versions/sb-5.0-version/use/release-notes/obsrv-2.0.0-ga)
{% endcontent-ref %}

{% content-ref url="/pages/8XuNwCmOdqszNFHwXIGj" %}
[Obsrv 2.2.0](/previous-versions/sb-5.0-version/use/release-notes/obsrv-2.2.0)
{% endcontent-ref %}

{% content-ref url="/pages/rkgbk5eBLIaEivoe8at4" %}
[Obsrv 2.1.0](/previous-versions/sb-5.0-version/use/release-notes/obsrv-2.1.0)
{% endcontent-ref %}

{% content-ref url="/pages/uX1IR13bqEVY4NH27Plz" %}
[Obsrv 5.3.0-GA](/previous-versions/sb-5.0-version/use/release-notes/release-v-5.3.0-ga)
{% endcontent-ref %}

{% content-ref url="/pages/EKmm4G0TKvcuA3ogJNnb" %}
[Obsrv 2.0-Beta](/previous-versions/sb-5.0-version/use/release-notes/release-v-5.2.0)
{% endcontent-ref %}

### Obsrv 1.0 Release Notes

{% content-ref url="/pages/EpLGphOVWJat3U2LP4wT" %}
[Release V 5.1.3](/previous-versions/sb-5.0-version/use/release-notes/release-v-5.1.3)
{% endcontent-ref %}

{% content-ref url="/pages/7aIhoFLuCN2fBzCtm5Ql" %}
[Release V 5.1.2](/previous-versions/sb-5.0-version/use/release-notes/release-v-5.1.2)
{% endcontent-ref %}

{% content-ref url="/pages/WJTfjJjC2au9mlFUZ4qV" %}
[Release V 5.1.0](/previous-versions/sb-5.0-version/use/release-notes/release-v-5.1.0)
{% endcontent-ref %}

{% content-ref url="/pages/1OD0v4QusT86ZkZPw0Ni" %}
[Release V 5.0.0](/previous-versions/sb-5.0-version/use/release-notes/release-v-5.0.0)
{% endcontent-ref %}

{% content-ref url="/pages/TBL1DwjgSC1Pi356GBie" %}
[Release V 4.10.0](/previous-versions/sb-5.0-version/use/release-notes/release-v-4.10.5)
{% endcontent-ref %}


# Obsrv 2.0-Beta

## **Features**

#### **Functional**

* \[Design] Generalisation of Obsrv core with domain agnostic [#OB-344](https://project-sunbird.atlassian.net/browse/OB-344)
* \[Design] Datasource Create API [#OB-338](https://project-sunbird.atlassian.net/browse/OB-338)
* \[Design] Datasource List API [#OB-332](https://project-sunbird.atlassian.net/browse/OB-332)
* \[Design] Dataset List API [#OB-314](https://project-sunbird.atlassian.net/browse/OB-314)
* \[Design] Dataset Update API [#OB-312](https://project-sunbird.atlassian.net/browse/OB-312)
* \[Design ] Dataset create API [#OB-306](https://project-sunbird.atlassian.net/browse/OB-306)
* \[Design] Data In API [#OB-300](https://project-sunbird.atlassian.net/browse/OB-300)
* \[Implementation] Data In API [#OB-288](https://project-sunbird.atlassian.net/browse/OB-288)
* Design for configuration of telemetry data specification for a data source [#OB-112](https://project-sunbird.atlassian.net/browse/OB-112)
* API implementation for Telemetry data specification configuration for a data source [#OB-118](https://project-sunbird.atlassian.net/browse/OB-118)
* \[Implementation] Datasource Update API [#OB-362](https://project-sunbird.atlassian.net/browse/OB-362)
* \[Implementation] Datasource Create API [#OB-356](https://project-sunbird.atlassian.net/browse/OB-356)
* \[Implementation] Dataset List API [#OB-320](https://project-sunbird.atlassian.net/browse/OB-320)
* \[Implementation] Dataset Update API [#OB-313](https://project-sunbird.atlassian.net/browse/OB-313)
* \[Implementation] dataset create API [#OB-294](https://project-sunbird.atlassian.net/browse/OB-294)

#### **Devops**

* \[Deployment] One click installation of Obsrv on AWS using Terraform [#OB-100](https://project-sunbird.atlassian.net/browse/OB-100)
* \[Deployment] One click installation of Obsrv on Azure using Terraform [#OB-100](https://project-sunbird.atlassian.net/browse/OB-100)
* \[Deployment] Helm chart Creation for the APIs [#OB-386](https://project-sunbird.atlassian.net/browse/OB-386)
* \[Deployment] Docker image and Helm chart For Dataset APIs. [#OB-326](https://project-sunbird.atlassian.net/browse/OB-326)
* \[Deployment] Build and Deployment Automation Scripts [#OB-386](https://project-sunbird.atlassian.net/browse/OB-386)
* \[Deployment] Helm Chart For Query API [#OB-350](https://project-sunbird.atlassian.net/browse/OB-350)

#### **Quality**

* \[Unit Tests] Improving the Code Coverage for APIs [#OB-350](https://project-sunbird.atlassian.net/browse/OB-350)
* \[Unit Tests] Improving the Code Coverage for Obsrv Core [#OB-350](https://project-sunbird.atlassian.net/browse/OB-350)

#### **Documentations**

* \[Documentation] APIs Swagger Documentation Update [#OB-94](https://project-sunbird.atlassian.net/browse/OB-94)


# Obsrv 2.1.0

Release Date - 31st Aug'23

## **Features**

***

#### **Functional**

* \[Enhancement] Obsrv Query Wrapper APIs Implementation [#OB-543](https://project-sunbird.atlassian.net/browse/OB-543)
* \[Enhancement] Data Exhaust API Implementation [#OB-544](https://project-sunbird.atlassian.net/browse/OB-5442)
* \[Enhancement] Configure the Obsrv superset to use the Obsrv meta APIs [#OB-545](https://project-sunbird.atlassian.net/browse/OB-545)
* \[Enhancement] Datasource API enhancement to validate the Ingestion spec [#OB-546](https://project-sunbird.atlassian.net/browse/OB-546)
* \[Enhancement] Wrapper API Implementation to support the Ingestion spec submission [#OB-547](https://project-sunbird.atlassian.net/browse/OB-547)
* \[Enhancement] APIs Data IN/OUT Metrics Generation [#OB-548](https://project-sunbird.atlassian.net/browse/OB-548)

#### **Devops**

* \[Devops] Enabling the labels for all the services [#OB-549](https://project-sunbird.atlassian.net/browse/OB-549)
* \[Devops] Configure the obsrv to run with MinIO object store [#OB-550](https://project-sunbird.atlassian.net/browse/OB-550)
* \[Devops] Support on multi channel alerts [#OB-551](https://project-sunbird.atlassian.net/browse/OB-551)

#### **Documentations**

* \[Documentation] Open source Documentation Update On MinIO support [#OB-552](https://project-sunbird.atlassian.net/browse/OB-552)


# Obsrv 2.2.0

Release Date - 31st Oct'23

## **Features**

***

#### **Functional**

* \[Enhancement] Obsrv major vulnerabilties fixes [#OB-554](https://project-sunbird.atlassian.net/browse/OB-554)
* \[Enhancement] Bug fixes [#OB-555](https://project-sunbird.atlassian.net/browse/OB-555)


# Obsrv 2.0.0-GA

Release Date - 31st Dec'23

## **Features**

***

#### **Functional**

* \[Enhancement] **Enhancements in datasets management** [#OB-556](https://project-sunbird.atlassian.net/browse/OB-556) These feature enhances data management capabilities and provides flexibility in dataset maintenance.
  * Users can now delete draft datasets using the API.
  * User can now retire live datasets using the API.
  * Ability to configure the denorm on the master datasets
* \[Enhancement] **Core pipeline enhancments** [#OB-556](https://project-sunbird.atlassian.net/browse/OB-556) Core pipeline changes to route all the failed events with detailed summary to failed topic.
* \[Enhancement] **Detailed Debugging with Query Store** [#OB-556](https://project-sunbird.atlassian.net/browse/OB-556) Indexing all failed events into the query store for comprehensive and detailed debugging. This feature improves troubleshooting capabilities and accelerates issue resolution.
* \[Enhancement] **Automated Backup Configuration** [#OB-556](https://project-sunbird.atlassian.net/browse/OB-556) Automated the configuration of the backup system during installation, allowing users to define their preferred timezone. This ensures a seamless and customized backup setup for enhanced data protection.
* \[Enhancement] **Filtered Rollups** [#OB-556](https://project-sunbird.atlassian.net/browse/OB-556) Ability to create rollup on a filter using API
* \[Enhancement] **Bug Fixes and Improvements** [#OB-556](https://project-sunbird.atlassian.net/browse/OB-556)
  * Fixed the auto-conversion of the timestamp property to a string.
  * Automatically convert numeric field to Integer during ingestion into Query Store/Druid.
  * Improve the code coverage of obsrv core to 100%
* \[Feature] **Aggregated and Filtered Data Sources** [#OB-557](https://project-sunbird.atlassian.net/browse/OB-557) Introduced the option to create aggregated data sources and filtered data sources through the API. Also enhanced API to query on the rollup datasources. This feature provides users with more control over data sources, enabling customization based on specific requirements.
* \[Feature] **Connector Ecosystem** [#OB-557](https://project-sunbird.atlassian.net/browse/OB-557) These connectors enable the data IN and OUT of the platform and expand the reach of our platform.
  * JDBC connector - This connector supports popular databases such as MySQL, PostgreSQL.
  * Data Stream Source Connector - This connector supports real-time streaming data sources such as Apache Kafka.


# Obsrv 5.3.0-GA

## **Features**

#### **Functional**

* \[Design] Connectors Design Review [#OB-452](https://project-sunbird.atlassian.net/browse/OB-452)
* \[Implementation] Kafka Connector [#OB-458](https://project-sunbird.atlassian.net/browse/OB-458)
* \[Design] Denorm Connector [#OB-190](https://project-sunbird.atlassian.net/browse/OB-190)
* \[Implementation] Denorm Connector [#OB-464](https://project-sunbird.atlassian.net/browse/OB-464)
* \[Design] CRUD API's for Datasource Config and Dataset transformation [#OB-428](https://project-sunbird.atlassian.net/browse/OB-428)
* \[Implementation] CRUD API's for Datasource Config and Dataset transformation [#OB-476](https://project-sunbird.atlassian.net/browse/OB-476)
* \[Design] Update /dataset/save and /datasources/save api to incorporate changes for master datasets [#OB-422](https://project-sunbird.atlassian.net/browse/OB-422)
* \[Implementation] Update /dataset/save and /datasources/save api to incorporate changes for master datasets [#OB-482](https://project-sunbird.atlassian.net/browse/OB-482)

#### **Devops**

* \[Devops] Deployment of Obsrv 2.0 in Sunbird Dev Environment [#OB-404](https://project-sunbird.atlassian.net/browse/OB-404)
* \[Devops] New Dev environment set up for Obsrv 2.0 in Sunbird [#OB-392](https://project-sunbird.atlassian.net/browse/OB-392)
* \[Devops] Prometheus Alert manager installation on Kubernetes [#OB-424](https://project-sunbird.atlassian.net/browse/OB-434)
* \[Devops] Alert configurations for various Functional alerts [#OB-446](https://project-sunbird.atlassian.net/browse/OB-446)
* \[Devops] Github Actions (Build and Deployment Script for the Connectors) [#OB-488](https://project-sunbird.atlassian.net/browse/OB-488)
* \[Upgrade] Update Lern and Knowlg building blocks to use the Telemetry API instead of directly pushing data into Kafka [#OB-440](https://project-sunbird.atlassian.net/browse/OB-440)
* \[Upgrade] Update telemetry service API to internally use `/data/in` APIs [#OB-416](https://project-sunbird.atlassian.net/browse/OB-416)
* \[Automation]Set up scripts to create Sunbird Telemetry and Summary datasets and corresponding data sources [#OB-410](https://project-sunbird.atlassian.net/browse/OB-410)
* \[Devops] Integration of build and deployment using Jenkins for Sunbird deployments [#OB-398](https://project-sunbird.atlassian.net/browse/OB-398)

#### **Quality**

* \[Testing] End to End testing of obsrv 2.0 with multiple datasets [#OB-470](https://project-sunbird.atlassian.net/browse/OB-470)
* \[Unit Tests] Improving the code coverage for `/datasources/save` API(s) [#OB-500](https://project-sunbird.atlassian.net/browse/OB-500)
* \[Unit Tests] Improving the code coverage of Datasource Config and Dataset transformation APIs [#OB-506](https://project-sunbird.atlassian.net/browse/OB-506)
* \[Unit Tests] Kafka Connector Code Coverage Improvment [#OB-512"](https://project-sunbird.atlassian.net/browse/OB-512)
* \[Unit Tests] Denorm Connector Code Coverage Improvment [#OB-518"](https://project-sunbird.atlassian.net/browse/OB-518)

#### **Documentations**

* \[Documentation] APIs Swagger Documentation Update [#OB-94](https://project-sunbird.atlassian.net/browse/OB-94)
* \[Documentation] Obsrv 2.0 Migration Documentation [#OB-494](https://project-sunbird.atlassian.net/browse/OB-494)

#### Pefromance Benchamrk

* \[Benchamrk] Obsrv 2.0 Components Benchmark [#OB-494](https://project-sunbird.atlassian.net/browse/OB-494)


# Release V 5.1.0

#### <mark style="color:blue;">5.1.0</mark>

**Release Details**

| Phases                            | Start Date  | End Date    |
| --------------------------------- | ----------- | ----------- |
| **Planning Phases**               | 01-AUG-2022 | 12-AUG-2022 |
| **Design Discussion**             | 15-AUG-2022 | 26-AUG-2022 |
| **Sprint 1**                      | 29-AUG-2022 | 16-SEP-2022 |
| **Sprint 2**                      | 19-SEP-2022 | 07-OCT-2022 |
| **Regression Testing & Releases** | 10-OCT-2022 | 04-NOV-2022 |
| **Production Release**            | 04-NOV-2022 | 04-NOV-2022 |
| **Bug fixes and Support**         | 07-NOV-2022 | 11-NOV-2022 |

**Features**

**Sprint 1**

* Obsrv Infra single click installation in the local mode [#OB-27](https://project-sunbird.atlassian.net/browse/OB-27)
* CSP Migration plan analysis [#OB-45](https://project-sunbird.atlassian.net/browse/OB-45)

**Sprint 2**

* Moving of all the repos into the obsrv community and deployment script changes. [#OB-39](https://project-sunbird.atlassian.net/browse/OB-39)
* Obsrv infra single click installation bug fixes in the Kubernets cluster to setup the druid. [#OB-51](https://project-sunbird.atlassian.net/browse/OB-51)
* Obsrv Query Engine Deployment Script. [#OB-82](https://project-sunbird.atlassian.net/browse/OB-82)
* Obsrv Query Engine API Implementation [#OB-88](https://project-sunbird.atlassian.net/browse/OB-88)
* Merging of the release branches to main branch [#OB-244](https://project-sunbird.atlassian.net/browse/OB-244)
* Obsrv HdInsight Cluster CSP Variable generalisation changes [#OB-232](https://project-sunbird.atlassian.net/browse/OB-232)
* Obsrv Druid CSP Variable generalisation [#OB-226](https://project-sunbird.atlassian.net/browse/OB-226)
* Obsrv Secor Service CSP Variable Generalisation [#OB-220](https://project-sunbird.atlassian.net/browse/OB-220)
* Obsrv Report APIs CSP Variable Generalisation [#OB-214](https://project-sunbird.atlassian.net/browse/OB-214)
* Obsrv Pipeline CSP Variable Generalisation [#OB-208](https://project-sunbird.atlassian.net/browse/OB-208)
* Report API enhancement to generate the signed URL for the relative paths [#OB-202](https://project-sunbird.atlassian.net/browse/OB-202)
* Obsrv stand alone spark machine provision script CSP variable generalisation [#OB-238](https://project-sunbird.atlassian.net/browse/OB-238)

**Github Tag Details**

| Component                               | Build Tag                                                                                                           | Deploy Tag                                                                                                 |
| --------------------------------------- | ------------------------------------------------------------------------------------------------------------------- | ---------------------------------------------------------------------------------------------------------- |
| **Analytics Service**                   | [**release-5.1.0\_RC2**](https://github.com/Sunbird-Obsrv/sunbird-analytics-service/releases/tag/release-5.1.0_RC2) | [**release-5.1.0\_RC1**](https://github.com/project-sunbird/sunbird-devops/releases/tag/release-5.1.0_RC1) |
| **Spark Provision**                     | NA                                                                                                                  | [**release-5.1.0\_RC1**](https://github.com/project-sunbird/sunbird-devops/releases/tag/release-5.1.0_RC1) |
| **Spark HD Insights Cluster Provision** | NA                                                                                                                  | [**release-5.1.0\_RC1**](https://github.com/project-sunbird/sunbird-devops/releases/tag/release-5.1.0_RC1) |
| **Secor**                               | NA                                                                                                                  | [**release-5.1.0\_RC1**](https://github.com/project-sunbird/sunbird-devops/releases/tag/release-5.1.0_RC1) |
| **Flink Jobs**                          | NA                                                                                                                  | [**release-5.1.0\_RC1**](https://github.com/project-sunbird/sunbird-devops/releases/tag/release-5.1.0_RC1) |
| **Druid**                               | NA                                                                                                                  | [**release-5.1.0\_RC1**](https://github.com/project-sunbird/sunbird-devops/releases/tag/release-5.1.0_RC1) |
| **Sunbird Analytics Core**              | NA                                                                                                                  | [**release-5.1.0\_RC1**](https://github.com/project-sunbird/sunbird-devops/releases/tag/release-5.1.0_RC1) |

#### **Configurations**

| Service                        | Old Configurations                                                                                                                                                                                                                                                                                                                                                                          | New Configurations                                                                                                                                                                                                                                                                                       |
| ------------------------------ | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| **Analytics Service**          | azure\_private\_account\_name, azure\_private\_account\_secret, azure\_public\_account\_secret, azure\_public\_account\_name                                                                                                                                                                                                                                                                | cloud\_private\_storage\_accountname, cloud\_private\_storage\_secret, cloud\_public\_storage\_secret, cloud\_public\_storage\_accountname                                                                                                                                                               |
| **Spark Provision**            | sunbird\_private\_storage\_account\_name, sunbird\_private\_storage\_account\_key, sunbird\_public\_storage\_account\_name, sunbird\_public\_storage\_account\_key, s3\_storage\_key, s3\_storage\_secret, sunbird\_private\_azure\_report\_container\_name, sunbird\_public\_azure\_report\_container\_name, azure\_private\_storage\_account\_name, azure\_private\_storage\_account\_key | cloud\_private\_storage\_accountname, cloud\_private\_storage\_secret, cloud\_public\_storage\_accountname, cloud\_public\_storage\_secret, cloud\_storage\_privatereports\_bucketname, cloud\_storage\_publicreports\_bucketname, cloud\_private\_storage\_accountname, cloud\_private\_storage\_secret |
| **Spark HD Insight Provision** | sunbird\_private\_storage\_account\_name, sunbird\_private\_storage\_account\_key,                                                                                                                                                                                                                                                                                                          | cloud\_private\_storage\_accountname, cloud\_private\_storage\_secret                                                                                                                                                                                                                                    |
| **Secor**                      | secor\_azure\_container\_name, sunbird\_private\_storage\_account\_key, sunbird\_private\_storage\_account\_name, azure\_container\_name,                                                                                                                                                                                                                                                   | cloud\_storage\_telemetry\_bucketname, cloud\_private\_storage\_secret, cloud\_private\_storage\_accountname, cloud\_storage\_telemetry\_bucketname                                                                                                                                                      |
| **Flink**                      | sunbird\_private\_storage\_account\_name, sunbird\_private\_storage\_account\_key, checkpoint\_store\_type, sunbird\_private\_storage\_account\_name, sunbird\_private\_storage\_account\_key, s3\_access\_key, s3\_secret\_key, s3\_endpoint, s3\_path\_style\_access                                                                                                                      | cloud\_private\_storage\_accountname, cloud\_private\_storage\_secret, cloud\_service\_provider, cloud\_private\_storage\_endpoint, cloud\_storage\_pathstyle\_access, cloud\_private\_storage\_project                                                                                                  |
| **Druid**                      | sunbird\_private\_storage\_account\_name, sunbird\_private\_storage\_account\_key,sunbird\_druid\_storage\_account\_name, sunbird\_druid\_storage\_account\_key, druid\_azure\_container\_name, s3\_storage\_key, s3\_storage\_secret, s3\_storage\_container, s3\_storage\_endpoint,s3\_path\_style\_access, s3\_default\_bucket\_location,                                                | cloud\_private\_storage\_accountname, cloud\_private\_storage\_secret, cloud\_storage\_telemetry\_bucketname, cloud\_storage\_pathstyle\_access, cloud\_private\_storage\_project, cloud\_private\_storage\_endpoint,cloud\_private\_storage\_region, cloud\_storage\_telemetry\_type                    |
| **Sunbird Analytics Core**     | sunbird\_private\_storage\_account\_name, sunbird\_private\_storage\_account\_key, sunbird\_public\_storage\_account\_name, sunbird\_public\_storage\_account\_key                                                                                                                                                                                                                          | cloud\_private\_storage\_accountname, cloud\_private\_storage\_secret, cloud\_public\_storage\_accountname, cloud\_public\_storage\_secret, cloud\_storage\_telemetry\_type                                                                                                                              |


# Release V 5.1.2

#### <mark style="color:blue;">5.1.2 (Hot Fix)</mark>

**Release Details**

**Features**

**Sprint 1**

* CSP Refactor [#OB-525](https://project-sunbird.atlassian.net/browse/OB-525)

**Github Tag Details**

| Component              | Build Tag                                                                                                           | Deploy Tag         |
| ---------------------- | ------------------------------------------------------------------------------------------------------------------- | ------------------ |
| **Analytics Service**  | [**release-5.1.2\_RC1**](https://github.com/Sunbird-Obsrv/sunbird-analytics-service/releases/tag/release-5.1.2_RC2) | release-5.0.0\_RC1 |
|                        |                                                                                                                     |                    |
| **Analytics Core**     | [**release-5.1.2\_RC1**](https://github.com/Sunbird-Obsrv/sunbird-analytics-core/releases/tag/release-5.1.2_RC2)    | release-5.0.0\_RC1 |
| **Core Data Products** | [**release-5.1.2\_RC1**](https://github.com/Sunbird-Obsrv/sunbird-core-dataproducts/releases/tag/release-5.1.2_RC2) | release-5.0.0\_RC1 |

#### **Build Changes**

1. **Analytics Core Build Command**

> mvn clean install -DskipTests -DCLOUD\_STORE\_GROUP\_ID=org.sunbird -DCLOUD\_STORE\_ARTIFACT\_ID=cloud-store-sdk\_2.12 -DCLOUD\_STORE\_VERSION=1.4.0

2. **Analytics Core Data Products Build Command**

> mvn clean install -DskipTests -DCLOUD\_STORE\_GROUP\_ID=org.sunbird -DCLOUD\_STORE\_ARTIFACT\_ID=cloud-store-sdk\_2.12 -DCLOUD\_STORE\_VERSION=1.4.0

3. **Analytics Service Build Command**

> mvn clean install -DskipTests -DCLOUD\_STORE\_GROUP\_ID=org.sunbird -DCLOUD\_STORE\_ARTIFACT\_ID=cloud-store-sdk\_2.12 -DCLOUD\_STORE\_VERSION=1.4.0


# Release V 5.1.3

#### <mark style="color:blue;">5.1.3 (Hot Fix)</mark>

**Release Details**

**Features**

**Sprint 1**

* CSP Refactor [#OB-525](https://project-sunbird.atlassian.net/browse/OB-525)

**Github Tag Details**

| Component              | Build Tag                                                                                                           | Deploy Tag                                                                                                      |
| ---------------------- | ------------------------------------------------------------------------------------------------------------------- | --------------------------------------------------------------------------------------------------------------- |
| **Analytics Service**  | [**release-5.1.3\_RC1**](https://github.com/Sunbird-Obsrv/sunbird-analytics-service/releases/tag/release-5.1.3_RC1) | [**release-5.2.0\_RC6**](https://github.com/Sunbird-Obsrv/sunbird-data-pipeline/releases/tag/release-5.2.0_RC6) |
| **Spark Provision**    | NA                                                                                                                  | [**release-5.2.0\_RC6**](https://github.com/Sunbird-Obsrv/sunbird-data-pipeline/releases/tag/release-5.2.0_RC6) |
| **Analytics Core**     | [**release-5.1.3\_RC5**](https://github.com/Sunbird-Obsrv/sunbird-analytics-core/releases/tag/release-5.1.3_RC5)    | [**release-5.2.0\_RC6**](https://github.com/Sunbird-Obsrv/sunbird-data-pipeline/releases/tag/release-5.2.0_RC6) |
| **Core Data Products** | [**release-5.1.3\_RC5**](https://github.com/Sunbird-Obsrv/sunbird-core-dataproducts/releases/tag/release-5.1.3_RC5) | [**release-5.2.0\_RC6**](https://github.com/Sunbird-Obsrv/sunbird-data-pipeline/releases/tag/release-5.2.0_RC6) |
| **Secor**              | NA                                                                                                                  | [**release-5.2.0\_RC6**](https://github.com/Sunbird-Obsrv/sunbird-data-pipeline/releases/tag/release-5.2.0_RC6) |
| **Flink Jobs**         | [**release-5.2.0\_RC6**](https://github.com/Sunbird-Obsrv/sunbird-data-pipeline/releases/tag/release-5.2.0_RC6)     | [**release-5.2.0\_RC6**](https://github.com/Sunbird-Obsrv/sunbird-data-pipeline/releases/tag/release-5.2.0_RC6) |

#### **Build Changes**

1. **Analytics Core Build Command**

> mvn clean install -DskipTests -DCLOUD\_STORE\_GROUP\_ID=org.sunbird -DCLOUD\_STORE\_ARTIFACT\_ID=cloud-store-sdk\_2.12 -DCLOUD\_STORE\_VERSION=1.4.6

2. **Analytics Core Data Products Build Command**

> mvn clean install -DskipTests -DCLOUD\_STORE\_GROUP\_ID=org.sunbird -DCLOUD\_STORE\_ARTIFACT\_ID=cloud-store-sdk\_2.12 -DCLOUD\_STORE\_VERSION=1.4.6

3. **Analytics Service Build Command**

> mvn clean install -DskipTests -DCLOUD\_STORE\_GROUP\_ID=org.sunbird -DCLOUD\_STORE\_ARTIFACT\_ID=cloud-store-sdk\_2.12 -DCLOUD\_STORE\_VERSION=1.4.6


# Release V 5.0.0

| Component            | Build Tag                                                                                                      | Deploy Tag                                                                                                |
| -------------------- | -------------------------------------------------------------------------------------------------------------- | --------------------------------------------------------------------------------------------------------- |
| **CoreDataProducts** | [**release-5.0.0\_RC1**](https://github.com/project-sunbird/sunbird-core-dataproducts/tree/release-4.10.5_RC1) | [**release-5.0.0\_RC1**](https://github.com/project-sunbird/sunbird-data-pipeline/tree/release-5.0.0_RC1) |
| **EdDataProducts**   | [**release-5.0.0\_RC1**](https://github.com/Sunbird-Ed/sunbird-data-products/tree/release-5.0.0_RC1)           | [**release-5.0.0\_RC1**](https://github.com/project-sunbird/sunbird-data-pipeline/tree/release-5.0.0_RC1) |

## **Features**

* UCI Jobs Changes - Adding Session Identifier in the Report. [#SB-30518](https://project-sunbird.atlassian.net/browse/SB-30518)
* Cloud Storage SDK Changes. [#OB-5](https://project-sunbird.atlassian.net/browse/OB-5)
* HD Insights cluster - Spark compatible service from other CSP. [#OB-3](https://project-sunbird.atlassian.net/browse/OB-3)
* SECOR Compatibility check with GCP and OCI. [#OB-4](https://project-sunbird.atlassian.net/browse/OB-4)
* Sunbird Obsrv Single Click Installation. [#OB-1](https://project-sunbird.atlassian.net/browse/OB-1)
* Druid metadata & deep storage migration spike. [#OB-6](https://project-sunbird.atlassian.net/browse/OB-6)


# Release V 4.10.0

### <mark style="color:blue;">4.10.5</mark>

| Component                                             | Build Tag                                                                                                       | Deploy Tag                                                                                                  |
| ----------------------------------------------------- | --------------------------------------------------------------------------------------------------------------- | ----------------------------------------------------------------------------------------------------------- |
| <p><strong>Provision Jobs</strong><br>- Spark<br></p> | [**release-4.10.5\_RC1**](https://github.com/project-sunbird/sunbird-data-pipeline/tree/release-4.10.5_RC1)     |                                                                                                             |
| **AnalyticsCore**                                     | [**release-4.10.5\_RC1**](https://github.com/project-sunbird/sunbird-analytics-core/tree/release-4.10.5_RC1)    | [**release-4.10.5\_RC1**](https://github.com/project-sunbird/sunbird-data-pipeline/tree/release-4.10.5_RC1) |
| **CoreDataProducts**                                  | [**release-4.10.5\_RC1**](https://github.com/project-sunbird/sunbird-core-dataproducts/tree/release-4.10.5_RC1) | [**release-4.10.5\_RC1**](https://github.com/project-sunbird/sunbird-data-pipeline/tree/release-4.10.5_RC1) |
| **EdDataProducts**                                    | [**release-4.10.5\_RC1**](https://github.com/Sunbird-Ed/sunbird-data-products/tree/release-4.10.5_RC1)          | [**release-4.10.5\_RC1**](https://github.com/project-sunbird/sunbird-data-pipeline/tree/release-4.10.5_RC1) |

#### **Features**

* Java 11 Upgrade changes. [#SB-28238](https://project-sunbird.atlassian.net/browse/SB-28238)
* Spark 3.1 Upgrade changes. [#SB-28238](https://project-sunbird.atlassian.net/browse/SB-28238)


# Installation Guide

All components in Sunbird Obsrv can be installed through automation scripts. The automation scripts require [Ansible](https://docs.ansible.com/ansible/latest/index.html) as a prerequisite. Some of the components also require the de-facto package manager for Kubernetes, [helm](https://helm.sh/docs/) as a prerequisite to run the component on [Kubernetes](https://kubernetes.io).

#### Telemetry Service

The Sunbird Obsrv Telemetry service can be deployed onto Kubernetes using the helm chart. The deployment is handled using Ansible to manage the configuration and the commands that are necessary. The deployments can also be integrated into Jenkins, a popular CI/CD tool. The Telemetry Service deployment also has the capability of configurable horizontal scaling using the Horizontal Pod Scaling (HPA) concept of Kubernetes. A sample command to deploy the telemetry service on Kubernetes is provided below.

```
helm install telemetry-service sunbird-devops/kubernetes/helm_charts/telemetry 
-n <namespace> --create-namespace
```

{% embed url="<https://github.com/project-sunbird/sunbird-devops/tree/release-4.8.0/kubernetes/helm_charts/core/telemetry>" %}
Telemetry Service helm chart
{% endembed %}

#### Data Pipeline

Sunbird Obsrv Data Pipeline consists of a series of real-time streaming jobs chained together to unzip, transform and enrich the telemetry data. We use ansible and helm charts to deploy the series of jobs. The list of jobs that need to be deployed and their configurations can be controlled by the [ansible defaults configuration](https://github.com/project-sunbird/sunbird-data-pipeline/blob/release-4.8.0/kubernetes/ansible/roles/flink-jobs-deploy/defaults/main.yml#L168-L339).

```
ansible-playbook $currentWs/kubernetes/ansible/deploy_jobs.yml 
--extra-vars "chart_path=${currentWs}/kubernetes/helm_charts/datapipeline_jobs 
job_names_to_deploy=<comma-separate-list-of-job-names>"
```

{% embed url="<https://github.com/project-sunbird/sunbird-data-pipeline/tree/release-4.8.0/kubernetes/ansible/roles/flink-jobs-deploy>" %}
Ansible role for Data Pipeline
{% endembed %}

Sunbird Obsrv uses the following list of fields to de-normalize the user metadata. These fields are obtained by calling the user-read api belonging to [Sunbird Lern](https://lern.sunbird.org) building block. The de-normalization job can be modified to read the user metadata from a service/api of the adopter's choice.

```
firstName, lastName, encEmail, encPhone, language, rootOrgId, profileUserType (usertype, subusertype), 
userLocations(state, district, block, cluster, school), rootOrg (orgName), userId, 
framework, profileUserTypes (usertype, subusertype)
```

#### Data Service

The Data Service is a collection of data exhaust and report apis. The Data Service can be installed on Kubernetes using Ansible and Helm. A sample command to install the service is provided below. All the required configuration is managed using the [ansible configuration](https://github.com/project-sunbird/sunbird-devops/blob/release-4.8.0/ansible/roles/stack-sunbird/defaults/main.yml#L987-L1015).

```
helm install telemetry-service sunbird-devops/kubernetes/helm_charts/analytics 
-n <namespace> --create-namespace
```

{% embed url="<https://github.com/project-sunbird/sunbird-devops/tree/release-4.8.0/kubernetes/helm_charts/core/analytics>" %}
Data Service helm chart
{% endembed %}

#### Report Service

{% embed url="<https://github.com/project-sunbird/sunbird-devops/tree/release-4.8.0/kubernetes/helm_charts/core/report>" %}
Report Service helm chart
{% endembed %}

#### Summarisers

{% embed url="<https://github.com/project-sunbird/sunbird-data-pipeline/tree/release-4.8.0/ansible/roles/data-products-deploy>" %}
Ansbile role for Summarizer data products
{% endembed %}


# Obsrv 2.0 Installation Guide

TODO

Please click [here](https://github.com/Sunbird-Obsrv/obsrv-automation/blob/main/INSTALLATION.md) to read the installation documentation.


# Getting Started with Obsrv Deployment Using Helm

Sunbird Obsrv is a high-performance, cost-effective data stack with several components such as ingestion, querying, processing, backup, visualisation and monitoring. Obsrv 2.0 can be either installed using `Terraform` (Infrastructure as Code tool) or using `Helm` (Kubernets Package Manager).

## **Prerequisites**

Obsrv runs completely on a Kubernetes cluster. A completely functional Kubernetes cluster is expected for a seamless Obsrv installation.

### Hardware

Obsrv can support a volume of 5 million events per day with an average size of each event to be around 5 kb with the following specifications.

* Kubernetes version of 1.25 or greater
* Minimum of 16 cores of CPU
* Minimum of 64 GB of RAM
* PersistentVolume support in the Kubernetes cluster
* Support for LoadBalancer service to externally expose some of the Obsrv services. Popular implementations such as MetalLB or Traefik can be used to expose the services using external IPs.

### Software

#### Helm

```bash
curl -fsSL -o get_helm.sh https://raw.githubusercontent.com/helm/helm/main/scripts/get-helm-3
chmod 700 get_helm.sh
./get_helm.sh
```

#### Helm Dependencies

Run the following `helm repo add` command to download the required dependencies for running Obsrv.

```bash
# Example
helm repo add prometheus https://prometheus-community.github.io/helm-charts
```

* monitoring - `https://prometheus-community.github.io/helm-charts`
* redis - `https://charts.bitnami.com/bitnami`
* loki (version - 4.8.0 ) - `https://grafana.github.io/helm-charts`
* promtail (version - 6.9.3 ) - `https://grafana.github.io/helm-charts`
* velero (version - 3.1.6 ) - `https://vmware-tanzu.github.io/helm-charts`

#### Source Code

Clone the [obsrv-automation](https://github.com/Sunbird-Obsrv/obsrv-automation) github repository. The required list of helm charts to deploy Obsrv will be under the `terraform/modules/helm` directory.

```bash
git clone https://github.com/Sunbird-Obsrv/obsrv-automation.git
cd obsrv-automation/terraform/modules/helm
```

#### Resources and Services

Please be advised that the list of resources will be completely different for different cloud service providers.

**Common**

The following list of buckets/containers need to be created for different services to store the data. This is applicable to Object Storage such as `MinIO/Ceph` as well.

* `flink-checkpoints`
* `velero-backup`
* `obsrv`

**AWS**

1. IAM role with `AmazonS3FullAccess` policy. Services such as Api, Druid, Flink, Secor need to read and write access to S3 buckets.
2. Velero is a service which provides backups of the entire Obsrv cluster state through snapshots. Velero backup service needs a restricted user access to upload the snapshot state onto S3. The following IAM role policy needs to be attached to user created for velero backup. The access keys needs to be generated for the velero backup user as well.

   ```json
   {
   "Statement": [
       {
       "Action": [
           "ec2:DescribeVolumes",
           "ec2:DescribeSnapshots",
           "ec2:CreateTags",
           "ec2:CreateVolume",
           "ec2:CreateSnapshot",
           "ec2:DeleteSnapshot"
       ],
       "Effect": "Allow",
       "Resource": "*"
       },
       {
       "Action": [
           "s3:GetObject",
           "s3:DeleteObject",
           "s3:PutObject",
           "s3:AbortMultipartUpload",
           "s3:ListMultipartUploadParts"
       ],
       "Effect": "Allow",
       "Resource": [
           "arn:aws:s3:::<velero-s3-container-name>/*"
       ]
       },
       {
       "Action": [
           "s3:ListBucket"
       ],
       "Effect": "Allow",
       "Resource": [
           "arn:aws:s3:::<velero-s3-container-name>"
       ]
       }
   ],
   "Version": "2012-10-17"
   }
   ```
3. Serive Accounts: Service accounts enable access of the S3 object storage without the need for the access keys. If you prefer to use keys instead, you can skip the creation of service accounts. The list of service accounts needed

* Dataset API with the name `dataset-api-sa`
* Druid with the name `druid-raw-sa`
* Flink with the name `flink-sa`
* Secor with the name `secor-sa`

## **Deployment Instructions**

> Helm package manager provides an easy way to install specific components using a generic command. Configurations can be overriden by updating the `values.yaml` file in the respective Helm charts.

```bash
helm upgrade --install --atomic <release_name> <chart_name> -n <namespace> -f <path/values.yaml> --create-namespace --debug
```

### Prerequisites

#### Kubernetes Cluster Access

Helm package manager needs access to the Kubernetes cluster. The path to the KUBECONFIG file needs to be exported as an environment variable, either in the current shell or in environment configuration files such as `.bashrc`

```bash
export KUBECONFIG=<path_to_kubeconfig file>
```

### Postgres

Postgres is a RDBMS database which is used as the metadata store

```powershell
helm upgrade --install --atomic postgresql postgresql/postgresql-helm-chart -n postgresql --create-namespace --debug
```

### Redis

Redis is an in-memory key-value store primarily used as a distributed cache

```powershell
helm upgrade --install --atomic obsrv-redis redis/redis -n redis -f redis/values.yaml --create-namespace --debug
```

### Prometheus

Prometheus is a monitoring system with a dimensional data model, flexible query language, efficient time series database and modern alerting approach.

```powershell
helm upgrade --install --atomic monitoring monitoring/kube-prometheus-stack -n monitoring -f monitoring/values.yaml --create-namespace --debug
```

### Kafka

Apache Kafka is a distributed event store and stream-processing platform.

```powershell
helm upgrade --install --atomic kafka kafka/kafka-helm-chart -n kafka --create-namespace --debug
```

The following list of kafka topics are created by default. If you would like to add more topics to the list, you can do so by adding it to `provisioning.topics` configuration in the [values.yaml](https://github.com/Sunbird-Obsrv/obsrv-automation/blob/main/terraform/modules/helm/kafka/kafka-helm-chart/values.yaml) file.

* dev.ingest
* masterdata.ingest

### Druid

Druid is a high performance, real-time analytics database that delivers sub-second queries on streaming and batch data at scale

#### Druid CRD

```powershell
helm upgrade --install --atomic druid-operator druid_operator/druid-operator-helm-chart -n druid-raw --create-namespace --debug
```

#### Druid Cluster

Druid requires the following set of configurations to be provided for specific storage systems such as AWS S3, Azure Blob Storage, GCP Storage or MinIO/Ceph

**AWS**

```yaml
druid_deepstorage_type: s3
druid.extensions.loadList: ["druid-s3-extensions"]
# S3 Access keys
s3_access_key: ""
s3_secret_key: ""
s3_bucket: "obsrv"
```

**MinIO/Ceph**

```yaml
druid_deepstorage_type: s3
druid.extensions.loadList: ["druid-s3-extensions"]
# S3 Access keys
s3_access_key: ""
s3_secret_key: ""
# Use the ClusterIP of the MinIO service instead of the Kubernetes service name
# We have noticed that the service names don't resolve properly
druid_s3_endpoint_url: http://172.20.126.232:9000/
s3_bucket: "obsrv"
druid_s3_endpoint_signingRegion: "us-east-2"
```

**Azure**

```yaml
druid.extensions.loadList: ["druid-azure-extensions"]
druid_deepstorage_type: azure
azure_storage_account_name: ""
azure_storage_account_key: ""
azure_storage_container: "obsrv"
```

**GCP**

```yaml
druid.extensions.loadList: ["druid-google-extensions"]
druid_deepstorage_type: google
# Google cloud credentials json file where the access_token and credentials are stored.
google_application_credentials: 
gcs_bucket: "obsrv"
```

[**Hadoop**](https://druid.apache.org/docs/latest/development/extensions-core/hdfs/)

```yaml
druid_deepstorage_type: "hdfs"
# Include the "druid-hdfs-storage" extension as part of the existing the extensions list
druid.extensions.loadList: ["druid-hdfs-storage"]
druid.indexer.logs.directory: "/druid/indexing-logs"
druid.storage.storageDirectory: "/druid/segments"

```

```powershell
helm upgrade --install --atomic druid-raw druid_raw_cluster/druid-raw-cluster-helm-chart -n druid-raw --create-namespace --debug
```

### API

This service provides metadata APIs related to various resources such as datasets/datasources in Obsrv. The following configurations need to be specified in the [values.yaml](https://github.com/Sunbird-Obsrv/obsrv-automation/blob/main/terraform/modules/helm/dataset_api/dataset-api-helm-chart/values.yaml) file.

#### AWS

```yaml
exhaust_service.CONTAINER: obsrv
exhaust_service.CONTAINER_STORAGE_PROVIDER: aws
exhaust_service.CONTAINER_STORAGE_REGION: us-east-2
```

```powershell
helm upgrade --install --atomic dataset-api dataset_api/dataset-api-helm-chart -n dataset-api --create-namespace --debug 
```

### Flink Streaming Jobs

Flink jobs are used to process and enrich the data ingested into Obsrv in near-realtime.

#### Configuration Overrides

**AWS**

```yaml
checkpoint_store_type: s3
# S3 Access keys
s3_access_key: ""
s3_secret_key: ""
# Under base_config in the values.yaml
base.url: s3://flink-checkpoints
```

**MinIO/Ceph**

```yaml
checkpoint_store_type: s3
# S3 Access keys
s3_access_key: ""
s3_secret_key: ""
# Use the ClusterIP of the MinIO service instead of the Kubernetes service name
# We have noticed that the service names don't resolve properly
s3_endpoint:  http://172.20.126.232:9000/
# Under base_config in the values.yaml
base.url: s3://flink-checkpoints
```

**Azure**

```yaml
checkpoint_store_type: azure
azure_account: ""
azure_secret: ""
# Under base_config in the values.yaml
base.url: blob://flink-bucket
```

**GCP**

```yaml
checkpoint_store_type: gcp
# Google cloud credentials json file where the access_token and credentials are stored.
google_application_credentials: ""
base.url: blob://flink-bucket
```

[**Hadoop**](https://nightlies.apache.org/flink/flink-docs-master/docs/ops/state/checkpoints/#available-checkpoint-storage-options)

```yaml
checkpoint_store_type: hdfs
# Under base_config in the values.yaml
base.url: hdfs:///flink-bucket/
```

#### Flink Merged Pipeline Job

```powershell
helm upgrade --install --atomic merged-pipeline flink/flink-helm-chart -n flink --set image.registry=sunbird --set image.repository=sb-obsrv-merged-pipeline --create-namespace --debug
```

#### Flink Master Data Processor Job

```powershell
helm upgrade --install --atomic master-data-processor flink/flink-helm-chart -n flink --set image.registry=sunbird --set image.repository=sb-obsrv-master-data-processor --create-namespace --debug
```

### Backup Processes

#### Secor

**Configuration Overrides**

**AWS**

```yaml
# S3 upload manager which is responsible to upload backup to deepstorage.
upload_manager: com.pinterest.secor.uploader.S3UploadManager
cloud_store_provider: S3
aws_access_key: ""
aws_secret_key: ""
aws_region: us-east-2
```

**MinIO/Ceph**

```yaml
# S3 upload manager which is responsible to upload backup to deepstorage.
upload_manager: com.pinterest.secor.uploader.S3UploadManager
cloud_store_provider: S3
aws_access_key: ""
aws_secret_key: ""
# Use the ClusterIP of the MinIO service instead of the Kubernetes service name
# We have noticed that the service names don't resolve properly
aws_endpoint: http://172.20.126.232:9000/
aws_region: us-east-2
```

**Azure**

```yaml
upload_manager: com.pinterest.secor.uploader.AzureUploadManager
cloud_store_provider: Azure
azure_account_name: ""
azure_account_key: ""
```

**GCP**

```yaml
upload_manager: com.pinterest.secor.uploader.GsUploadManager
# Credentials path where access token and secrets are stored.
gs_credentials_path: google_app_credentials.json
```

**Hadoop**

```yaml
upload_manager: com.pinterest.secor.uploader.HadoopS3UploadManager
# Ensure the secor.s3.filesystem property is updated with the `hdfs` value
cloud_store_provider=hdfs
cloud_storage_bucket=namenode-host:8020/dir_path
# For More details please check here - https://github.com/pinterest/secor/issues/129
```

Secor backups are performed from various kafka topics which are part of the data processing pipeline. The following list of backup names need to be replaced in the below mentioned command.

List of backup names

* ingest-backup
* extractor-duplicate-backup
* extractor-failed-backup
* raw-backup
* failed-backup
* invalid-backup
* unique-backup
* duplicate-backup
* denorm-backup
* denorm-failed-backup
* system-stats
* system-events

```powershell
helm upgrade --install --atomic <backup_name> secor/secor-helm-chart -n secor --create-namespace
```

#### Velero

```powershell
helm upgrade --install --atomic velero velero/velero -n velero -f velero/values.yaml --create-namespace --debug --version 3.1.6
```

### Monitoring Services

#### Monitoring Dashboards

```powershell
helm upgrade --install --atomic grafana-configs grafana_configs/grafana-configs-helm-chart -n monitoring --create-namespace --debug
```

#### Monitoring Alert Rules

```powershell
helm upgrade --install --atomic alertrules alert_rules/alert-rules-helm-chart -n monitoring --create-namespace --debug
```

#### Druid Exporter

```powershell
helm upgrade --install --atomic druid-exporter druid_exporter/druid-exporter-helm-chart -n druid-raw --create-namespace --debug
```

#### Kafka Exporter

```powershell
helm upgrade --install --atomic kafka-exporter kafka_exporter/kafka-exporter-helm-chart -n kafka --create-namespace --debug
```

#### Postgres Exporter

```powershell
helm upgrade --install --atomic postgresql-exporter postgresql_exporter/postgresql-exporter-helm-chart -n postgresql --create-namespace --debug
```

#### Loki

```powershell
helm upgrade --install --atomic loki loki/loki -n loki -f loki/values.yaml --create-namespace --debug --version 4.8.0
```

#### Promtail

```powershell
helm upgrade --install --atomic promtail promtail/promtail -n loki -f promtail/values.yaml --create-namespace --debug --version 6.9.3
```

### Ingestion

This helm chart is used to submit the default ingestion tasks required for the system statistics events

```powershell
helm upgrade --install --atomic submit-ingestion submit_ingestion/submit-ingestion-helm-chart -n submit-ingestion --create-namespace --debug 
```

### Visualization

#### Superset

```powershell
helm upgrade --install --atomic superset superset/superset-helm-chart -n superset --create-namespace --debug
```

### Loadbalancers

Following is a list of services which are exposed as a LoadBalancer service.

| Component   | Service Name                | Description             |
| ----------- | --------------------------- | ----------------------- |
| Dataset API | service/dataset-api-service | Meta APIs               |
| Superset    | service/superset            | Data Visualization Tool |

## Post Deployment

Please find documentation related to various application level functionalities in Obsrv below

* [Create Datasets](https://github.com/Sunbird-Obsrv/obsrv-automation/blob/main/INSTALLATION.md#create-a-dataset)
* [Data Ingestion](https://github.com/Sunbird-Obsrv/obsrv-automation/blob/main/INSTALLATION.md#data-ingestion)
* [Data Querying](https://github.com/Sunbird-Obsrv/obsrv-automation/blob/main/INSTALLATION.md#data-query)


# System Requirements

### Requirements for Telemetry Service

| Software | Version       |
| -------- | ------------- |
| Node     | 12x or above  |
| Docker   | 19.x or above |
| Kafka    | 2.4.1         |

### Requirements for Data Pipeline

| Software     | Version         |
| ------------ | --------------- |
| Java         | 11              |
| Scala        | 2.12.x          |
| Apache Flink | 1.13.0          |
| Kafka        | 2.4.1           |
| Redis        | 5.x or above    |
| Kubernetes   | 1.17.0 or above |

### Requirements for Data Service

| Software | Version       |
| -------- | ------------- |
| Java     | 8             |
| Scala    | 2.11.x        |
| Docker   | 19.x or above |
| Postgres | 9.6 or above  |

### Requirements for Report Service

| Software | Version       |
| -------- | ------------- |
| Node     | 12x or above  |
| Docker   | 19.x or above |
| Postgres | 9.6 or above  |

### Requirements for Report Configurator

| Software | Version       |
| -------- | ------------- |
| python   | 3.5 or above  |
| Docker   | 19.x or above |

### Requirements for Summarisers

| Software | Version |
| -------- | ------- |
| Java     | 8       |
| Scala    | 2.11.x  |
| Spark    | 2.4.4   |
| Kafka    | 2.4.1   |

###


# LEARN


# Functional Capabilities

Functional capabilities that can be enabled using Sunbird Obsrv

The following are the key functional capabilities that one would turn to Sunbird Obsrv for:

**Flexible Data Model Specification to adapt to any domain:** The Telemetry data model is designed to instrument telemetry data (interactions, metrics and logs) for any domain.

The domain of an application determines to a vast extent the design of the data model for understanding various interactions and metrics from the application. Sunbird Obsrv is flexible to allows the application to decide on the data model for the Telemetry and it is designed intelligently to keep the data structure for the event data as flexible as possible while enforcing to capture some common characteristics to ascertain the component details of the source system that generates the Telemetry. The flexible data model also ensures that the downstream processing systems have to change very little to accommodate the changes in the data model and allows re-usability of the majority of the components in the analytics platform.

The data model works based on the telemetry specifications defined as per the domain of application, based on the [Sunbird Telemetry](https://telemetry.sunbird.org) building block.

**Data collection at scale:** Perform seamless telemetry data collection from client apps using horizontally scalable apis<mark style="color:green;">.</mark>

The ability to seamlessly scale the telemetry data collection is very important as it enables the client applications to sync the data without any data loss and also determines the rate at which . The data collection api in Sunbird Obsrv can be horizontally scaled by adding more instances of the collector both in containerized and virtual machine based environments. The data collector also provides options for configurable data sink such as Apache Kafka, Apache Cassandra or local hard disks. The data sinks, in turn, can be scaled horizontally which helps transferring telemetry data at very high transfer rates along with consistent and durable storage.

This is enabled by the [Telemetry Service](/previous-versions/sb-5.0-version/learn/product-and-developer-guide/telemetry-service) component of the Obsrv building block

**Easy integration with Cloud providers:** The analytics platform is cloud agnostic and can be easily deployed onto any of the popular cloud platforms such as AWS, Azure or Google Cloud Platform.

The analytics platform has various plugable components and all of the components provide configurations to support deployments onto the popular cloud platforms. Furthermore, most of the components are deployed onto Kubernetes, the most popular tool for container orchestration.

**Streaming/Batch data processing:** The analytics platform provides an ability to process both real-time streaming data and batch data to derive insights.

The analytics platform is built using lambda architecture principles. Lambda architecture is designed to provide the capability to process both real-time streaming data and batch data. The analytics platform has both a batch layer and a stream layer storage which helps in fault-tolerant and scalable architecture for data processing. The stream layer ensures that insights are derived in near real time to take immediate actions on the data arriving and the batch layer ensures that reporting requirements over a specific time period are addressed. Fault tolerance is built into the pipeline which ensures that there will be no data loss. Lambda architecture also provides the capability to replay or reload data for data processing with idempotency.

The [Data Pipeline](/previous-versions/sb-5.0-version/learn/product-and-developer-guide/data-pipeline) component of the building block is used to enable these capabilities

**Build fast, scalable data analytics:** The analytics platform is designed to build ad-hoc analytics or recurring reports in a fast, scalable manner.

Any analytics platform should allow ad-hoc user queries to be served with low latency response times and should provide an immutable data storage to obtain a snapshot of the historical data. Sunbird Obsrv's analytics platform uses different data stores to address the ability to query data with low latency. It uses the distributed cloud storage for batch processing of data to generate reports and it uses an Online Analytical Processing (OLAP) data store to address real-time adhoc user queries. Both these layers provide the ability to build fast and scalable analytics and reporting systems on top of large amounts (terabytes) of data with efficient infrastructure. The access of the data from the polyglot data stores is always using APIs which can horizontally scale according to the latency requirements.

**Readymade metrics & Downloadable data:** The analytics platform provides an ability to download both the raw data and aggregated data through the reporting api also with an ability to schedule recurring reporting requirements.

Apart from the ability to build a querying layer on top of the data, a data analytics platform should provide the ability to download the data from the system so that the client applications can build their own insights on top of the data. Sunbird Obsrv's analytics platform is designed with this capability in mind. Cloud data stores in general store objects that are private and only the object owner has permissions to access them. The analytics platform provides APIs with access control to download data. The APIs generate urls which are pre-signed with security credentials of the object owner to grant time-bound permissions to download the data. This capability opens up a lot of possibilities for the client applications to create a visualization without overloading the data platform.

The quick-analysis capabilities, along with the readymade metrics and downloadable datasets are made available using the [Data Service](/previous-versions/sb-5.0-version/learn/product-and-developer-guide/data-service), [Report Service](/previous-versions/sb-5.0-version/learn/product-and-developer-guide/report-service), [Report Configurator](/previous-versions/sb-5.0-version/learn/product-and-developer-guide/report-configurator) and the [Summarisers](/previous-versions/sb-5.0-version/learn/product-and-developer-guide/summarisers). These components come together to enable a boquet of reporting services for the Sunbird Obsrv building block.

Some additional possibiities that can be visualised using Sunbird Obsrv include:

* Allowing for `capture of telemetry` as per chosen specifications from an agricultural monitoring system, and `processing and analysis of the resulting data` to optimise output
* Aggregation of anonymised health data from across platforms, and its `analysis` to arrive at health indicators and patterns at population scale.
* Analysis of `streams of weather data` from satellites to better `predict` weather patterns.

*Obsrv* allows for extensive configurability so as to enable creation of various custom workflows. Refer to the pages for different components listed in the Product & developer guide page to see details of what configurations are available for each.


# Dependencies

## Sunbird Telemetry <a href="#sunbird-telemetry" id="sunbird-telemetry"></a>

Sunbird Telemetry is a specification to instrument all the key events. Using this specification reference applications & services will generate telemetry events.

{% hint style="info" %}
Resolution : *Obsrv* as a building block serves to collect and process telemetry. Just like it works with the Sunbird telemetry spec, it can work seamlessly with any spec that a system may adopt as its choice of reference. The events generated will be in accordance with the spec of choice, and *Obsrv* can carry out validation and aggregate telemetry accordingly.
{% endhint %}

## Sunbird Lern

The User and Org service of Sunbird Lern is used for obtaining various metadata about the user - including fetching user roles & privileges to authenticate exhaust requests.

{% hint style="info" %}
Resolution: *Obsrv* uses the meta data provided by the *Lern* building block in order to gather user data. An alternate service can be used in place of L\_ern\_ to make this data available to *Obsrv* as well.
{% endhint %}

## Keycloak & Kong <a href="#sunbird-telemetry" id="sunbird-telemetry"></a>

\<description>

{% hint style="info" %}
Resolution :
{% endhint %}


# Product Roadmap

Sunbird Obsrv 2.0 has been redesigned from ground up to allow ingestion, processing and querying of the telemetry data to be agnostic of the data specification of the telemetry data. The last stable release of Sunbird Obsrv 1.0, which is tightly coupled of the [Sunbird Telemetry](https://telemetry.sunbird.org/) Specification will be [Release 5.1.1](/previous-versions/sb-5.0-version/use/release-notes/release-v-5.1.0). The product roadmap for Sunbird Obsrv 2.0 has been detailed out below with the features being logically grouped under specific functional features.

Sunbird Obsrv [ISSUE TRACKER](https://github.com/Sunbird-Obsrv/Community/issues) : This is the link to the set of issues/ submissions or requests that are being considered for development as part of the Sunbird Obsrv roadmap. You can upvote an issue if you find it relevant, or <mark style="color:blue;">add a new issue</mark> to the list

<mark style="color:orange;">**Obsrv AMJ-2024 Stories**</mark>

1. **APIs Refactoring**
   * Dataset CRUD API
   * Dataset Management API
   * Data In Out APIs
   * Query Template APIs
2. **Automation Refactoring**
   * Helm charts refactoring
   * Wrapper over helm charts
3. **Connectors**
   * Connectors Framework
   * Connectors Management
   * Connectors Implementation
4. **Hudi Integration**
   * Create an ingestion spec for Lakehouse tables
   * Streaming job implementation to write Raw Data to Lakehouse
   * Timestamp Based Partitioner for Apache Hudi
   * Hudi Sink Configuration Optimization
   * Deployment Automation for Hudi Sink Connector - AWS/Local Datacenter
   * Query API Unification for the Lakehouse/Real-Time store - Design
   * Query API Unification for the Lakehouse/Real-Time store - Implementation
   * Dedup Challenges with introduction of Lakehouse - Design
   * Dedup Challenges with introduction of Lakehouse - Implementation
   * Design Rollups for the Lakehouse
   * Rollups implementation for the Lakehouse
   * Deployment Automation for Hudi Sink Connector - Azure
   * Deployment Automation for Hudi Sink Connector - GCP

<mark style="color:orange;">**Obsrv 2.0.1 GA Release date - 29th Feb'24**</mark>

<mark style="color:orange;">**Enhancements**</mark>

1. **Core pipeline enhancments**
   * To handle denormalization logic for empty keys, text & numeric keys.
   * Moved failed events sinking into a common base class.
   * Updated framework to created dynamicKafkaSink object.
   * Master dataset processor can now do denormalization with another master dataset as well.
2. **Tech software upgrades** - Postgres, Druid, Superset, NodeJS, Kubernetes.
3. **Dataset Management** - Added timezone handling to store the data in druid in the TZ specified by the dataset.
4. **Infra reliability** - These feature enhances data management capabilities and provides flexibility in dataset maintenance.
   * Verification of metrics for various services.
   * Metrics intrumentation for few of the services.
5. **Enhancments to dataset API service** - API endpoint changes.
6. **Benchmarking** - Processing, Ingestion & Querying.
7. **Bug Fixes and Improvements**
   * Fix unit tests and improve coverage.
   * Automation script enhancements to support new changes.

<mark style="color:orange;">**Feature**</mark>

1. **Connector Framework implementation**
2. **Hudi Data Ingestion**

<mark style="color:orange;">**Obsrv 2.0.0 GA Release date - 31st Dec'23**</mark>

<mark style="color:orange;">**360 degree observability**</mark>

1. **Enhancements in datasets management** - These feature enhances data management capabilities and provides flexibility in dataset maintenance.
   * Users can now delete draft datasets using the API.
   * User can now retire live datasets using the API.
   * Ability to configure the denorm on the master datasets
2. **Core pipeline enhancments** - Core pipeline changes to route all the failed events with detailed summary to failed topic.
3. **Detailed Debugging with Query Store** - Indexing all failed events into the query store for comprehensive and detailed debugging. This feature improves troubleshooting capabilities and accelerates issue resolution.
4. **Automated Backup Configuration** - Automated the configuration of the backup system during installation, allowing users to define their preferred timezone. This ensures a seamless and customized backup setup for enhanced data protection.
5. **Aggregated and Filtered Data Sources** - Introduced the option to create aggregated data sources and filtered data sources through the API. Also enhanced API to query on the rollup datasources. This feature provides users with more control over data sources, enabling customization based on specific requirements.
6. **Filtered Rollups** - Ability to create rollup on a filter using API
7. **Bug Fixes and Improvements**
   * Fixed the auto-conversion of the timestamp property to a string.
   * Automatically convert numeric field to Integer during ingestion into Query Store/Druid.
   * Improve the code coverage of obsrv core to 100%

<mark style="color:orange;">**Connector Ecosystem**</mark>

These connectors enable the data IN and OUT of the platform and expand the reach of our platform.

1. **JDBC connector** - This connector supports popular databases such as MySQL, PostgreSQL.
2. **Data Stream Source Connector** - This connector supports real-time streaming data sources such as Apache Kafka.

<mark style="color:orange;">**Obsrv 2.2.0 Release date - 31st Oct'23**</mark>

<mark style="color:orange;">**360 degree observability**</mark>

1. Addressing major vulnerabilities, making Obsrv BB free from vulnerabilities.
2. Restarting Flink to pick up new datasets. The command service with the Flink restart command needs to be open-sourced.
3. Few Bug fixes with components & deployment.

<mark style="color:orange;">**Obsrv 2.1.0 Release date - 30th Sep'23**</mark>

<mark style="color:orange;">**360 degree observability**</mark>

1. Restarting Flink to pick up new datasets. The command service with the Flink restart command needs to be open-sourced.
2. Refactor of Sunbird Ed Cache Updater Jobs
3. Enhance the Device register/profile API to send transaction event to obsrv 2.0 system using data IN API or Kafka connector

<mark style="color:orange;">**Connector Ecosystem**</mark>

1. Object Storage Connectors - Cloud storages
2. MinIO connector
3. Connector Marketplace/Framework
4. Postgresql Connector for data & denormalization

<mark style="color:orange;">**Obsrv 2.1.0 Release date - 31st Aug'23**</mark>

<mark style="color:orange;">**360 degree observability**</mark>

1. Data Exhaust API - Ability to download the raw data using the Data Exhaust APIs
2. Druid Query and Ingestion Submission Wrapper APIs - Ability to query and ingestion spec submission wrapper API
3. APIs System to generate the Data IN/OUT Metrics Generation
4. Integrate the Obsrv superset with Query Wrapper APIs

<mark style="color:orange;">**Simplified Operations**</mark>

1. Configure the obsrv to run with MinIO object store
2. Support on multi channel alerts
3. Enabling the labels for all the services

[<mark style="color:orange;">**Obsrv 2.0.0 GA**</mark>](/previous-versions/sb-5.0-version/use/release-notes/release-v-5.3.0-ga)<mark style="color:orange;">**(Planned release date - 2 May'23)**</mark>

<mark style="color:orange;">**360 degree observability**</mark>

1. Data Set Creation via APIs
2. Data Denormalization with API & Push to Kafka
3. Real time querying
4. Druid SQL & JSON query interface
5. Out-of-the-box visualizations with Superset

<mark style="color:orange;">**Simplified Operations**</mark>

1. One-click installation for AWS, Azure & GCP
2. Standard Monitoring Capability
3. Log Streaming with Grafana UI

<mark style="color:orange;">**Connector Ecosystem**</mark>

1. Kafka connector for data & denormalization

<mark style="color:orange;">**Unified Web Console**</mark>

1. General Cluster metrics

\
[<mark style="color:orange;">**Release-5.1.0**</mark>](/previous-versions/sb-5.0-version/use/release-notes/release-v-5.1.0) <mark style="color:orange;">**(Planned release date - 04 Nov'22) - Projects**</mark>

<mark style="color:orange;">**Project : Enabling ease of adoption**</mark>

<mark style="color:orange;">**Task : One click install enhancements**</mark>

a. Include the data products (work flow summary) as part of the one-click installer package. These allow for easy generation of common summaries once raw telemetry is generated by the adopter.

b. Include a sample data set and some pre-configured charts as part of the one-click installer package so that an adopter trying it out can get a feel of the types of charts that can be generated, and the kind of data required to be generated for this purpose.

<mark style="color:orange;">**Project: Enabling ease of adoption**</mark>

<mark style="color:orange;">**Task : Additional documentation**</mark>

Additional pending Sunbird Obsrv documentation work

<mark style="color:orange;">**Project: Enabling ease of adoption**</mark>

<mark style="color:orange;">**Task: Creation of a Learning module for Sunbird Obsrv**</mark>

Compile existing learning resources/ create new resources in order to put together a learning module/ course that will allow a developer to get familiar with the Obsrv building block, and potentially get "Obsrv certified" by earning a certificate after taking an assessment.

<mark style="color:green;">**4.10.0 :**</mark> Planned for 06 Jun '22

| <p><strong>1. Functional Requirements: </strong><mark style="color:green;"><strong>4.10.0</strong></mark></p><p><strong>-</strong> API level accsess to data files that power the reports and charts on the portal</p>                                                                                                                                                                                                                                                                                                                                       |
| ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ |
| - Ability to configure Sunbird datasets as Public or Private at an instance level                                                                                                                                                                                                                                                                                                                                                                                                                                                                            |
| - KT for indexing variables into Druid                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                       |
| <p><strong>2. Deployment and Release Processes: </strong><mark style="color:green;"><strong>4.10.0</strong></mark></p><p>Build, Deploy and provisioning scripts : refactoring of the provisioning and deployment scripts for the BB- Sunbird dev and staging environments will be repurposed for each BB. The deployment scripts need to be refactored to sandbox the environment on Kubernetes for Sunbird Obsrv BB. This will be done by deploying the services and components on Kubernetes onto a separate configurable namespace for Sunbird BB<br></p> |
|                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                              |
| <mark style="color:green;">**5.0.0 :**</mark> Planned for 19 Aug'22                                                                                                                                                                                                                                                                                                                                                                                                                                                                                          |
| <p><strong>1. One click mini installation of data pipeline components on Kubernetes: </strong><mark style="color:green;"><strong>5.0.0</strong></mark></p><p><strong>-</strong> Currently not all components are deployed onto Kubernetes and it takes significant effort from the adopter to get all the components up and running. This capability will allow the adopter to quickly install required components on Kubernetes and have the entire analytics platform up and running.</p>                                                                  |
| <p><strong>2. Publish the minimum number of Kubernetes nodes required for one click installation: </strong><mark style="color:green;"><strong>5.0.0</strong></mark></p><p>- Analyse and publish minimum number of nodes required to get the analytics platform up and running on Kubernetes. Also, publish the mandatory components required to get the installation up and running.</p>                                                                                                                                                                     |
| <p><strong>3. Migration of API Swagger documentation to Sunbird Obsrv Building block pages: </strong><mark style="color:green;"><strong>5.0.0</strong></mark></p><p>- Migrate existing Sunbrid Swagger API for analytics API and data exhaust APIs documentation to Sunbird Obsrv Building block pages.</p>                                                                                                                                                                                                                                                  |
| <p><strong>4. Multi-cloud support for blob store:</strong> <mark style="color:green;"><strong>5.0.0</strong></mark></p><p>Update the framework, cloud-storage sdk, Secor and Flink checkpoints to work with GCP as well- Generalize the analytics framework, cloud-storage-sdk, Secor and Flink checkpoints to work with AWS S3, Azure Blob Storage and Google Cloud Storage- Add relevant configuration files required for the generalization<br></p>                                                                                                       |
| - Be able to deploy existing microservices into a different namespace (SB Ed)                                                                                                                                                                                                                                                                                                                                                                                                                                                                                |
|                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                              |
| [**Sunbird Obsrv Backlog >**](https://project-sunbird.atlassian.net/issues/?filter=12538)                                                                                                                                                                                                                                                                                                                                                                                                                                                                    |
| <p><strong>1. Documentation to configure data exhaust reports using the APIs:</strong> Detailed explanation of the different sections of the configuration schema.</p><p>- Add documentation to explain the configuration schema for the various data exhaust report API</p><p>- Add detailed documentation on the usage of the APIs</p>                                                                                                                                                                                                                     |
| <p><br></p>                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                  |


# Product & Developer Guide

**OVERVIEW**

Sunbird allows for generation and consumption of usage data in order to be able to carry out any required analysis. The data flow on the platform is as illustrated below.

The client apps and micro services within Sunbird are instrumented to generate Telemetry data to capture usage and various actions to derive insights. The data pipeline will process the Telemetry event stream data in near real-time by application of a series of validation, de-duplication, transformation and denormalization steps. The Telemetry data is synced at various points in the data pipeline to a configurable cloud storage which acts as the persistent data storage.

A generic summarized dataset is generated from the Telemetry data which will comprise of commonly used summaries such as session summary. The summarized dataset is an aggregated dataset which can be used to generate custom derived datasets required for custom workflows.

The telemetry event data is then synced to a configurable analytic data store to build a access layer to the data sets. The adopters are free to choose any analytics data store and/or visualisation tools (Sunbird supports Druid and Superset out-of-the-box) of their choice to build canned reports and dashboards to slice and dice the data.

Sunbird Obsrv is a combination of various tools which provide the capabilities such as streaming, processing and storage of telemetry data and deriving reporting insights from the data.

![L0 Sunbird Obsrv Component View](/files/6LhbYQoZxqouuAitBZTy)

![High Level Sunbird Obsrv Component Overview](/files/AWHqioZY2ZVMDCC94uNZ)

![Data analytics architecture](/files/5eUcZoDQSz5HX8yYT0A3)


# Telemetry Service

Telemetry Service is a highly scalable micro service which allows various components in the Sunbird platform to send the event level data generated from various activities in the platform. These event level data could be either from clients such as the user apps or from micro services such as Sunbird Lern.

![Telemetry Service Workflow](/files/uvRcNTc1oT0VmzGSOEZK)

#### Key Features:

1. Horizontally Scalable: The service can be scaled horizontally to handle increase in Telemetry data traffic as the operations are idempotent.
2. Configurable storage: The service supports disk storage, Apache Kafka, Apache Cassandra as data storage.
3. In-built compression: Efficient use of network bandwidth by compressing the messages sent to the data storage.

#### Installation Configuration Reference:

| Property            | Description                                                                                         | Default                               |
| ------------------- | --------------------------------------------------------------------------------------------------- | ------------------------------------- |
| localStorageEnabled | A boolean flag to enable on disk storage                                                            | true                                  |
| dispatcher          | Supported dispatchers are kafka, cassandra, file and console                                        | console                               |
| kafkaHost           | Host or IP address of the Kafka broker servers                                                      | none                                  |
| topic               | The kafka topic to which telemetry data will be pushed to                                           | none                                  |
| compression\_type   | Type of compression to use if kafka dispatcher is chosen. Possible values are none, gzip and snappy | none                                  |
| filename            | The name of the file to write data to if file dispatcher is chosen                                  | telemetry-%DATE%.log                  |
| maxSize             | Roll over max size for file based dispatched                                                        | 100 Mb                                |
| maxFiles            | Max files to keep for a file based dispatcher                                                       | 100                                   |
| partitionBy         | Time based partitioning key for a Cassandra dispatcher                                              | hour                                  |
| keyspace            | Keyspace name for a Cassandra dispatcher                                                            | none                                  |
| cassandraTtl        | TTL for retention of data for a Cassandra dispatcher                                                | none                                  |
| threads             | Number of cpu threads to be used by the service                                                     | Default number of cpus on the machine |

{% hint style="info" %}
Telemetry documentation [telemetry.sunbird.org](http://127.0.0.1:5000/o/-Mi9QwJlsfb7xuxTBc0J/s/-MkM7F4oILSpCJPO0YUu/)
{% endhint %}

{% embed url="<https://github.com/project-sunbird/sunbird-telemetry-service>" %}
Source Code
{% endembed %}


# Data Pipeline

Group of real-time stream processing jobs that processes the event stream of telemetry data generated from client apps and micro services. The telemetry data goes through a series of steps such as validation, de-duplication, transformation and denormalization of metadata. The transformed data is then stored in a consumable format that can be used for further analysis.

![Analytics Data Pipeline](/files/n7avIQnOsi4G0uy13fNb)

#### Key Features:

1. Lambda Architecture: A hybrid approach of using both batch-processing and stream-processing methods to process massive data sets.
2. Loose coupling: Data processing jobs are loosely coupled as they only communicate with a durable queue such as Apache Kafka.
3. Easy chaining: Data processing jobs can be chained easily by only configuring the input and output data sources they consume from. This allows easy introduction of new jobs required for processing custom workflows.
4. Data Sync points: The stream of data is synced to a configurable cloud storage which acts as a persistent data store with durability. The data sync points allow the capability to replay data from a specific stage in the pipeline.
5. Resiliency: The data pipeline guarantees AT LEAST ONCE processing semantics and ensures no data loss.
6. Monitoring: The data pipeline jobs has the capability to emit standard and custom metrics to monitor the health and also allows you to perform an audit of the system at various stages.
7. Auto-scaling: The data processing pipeline offers support for auto-scaling out of the box. This helps the pipeline to adapt to changes in the incoming data volume.
8. Real-time analytics: The data processing pipeline offers out of the box support with Apache Druid, an analytics data store design for fast slice-and-dice analytics.

#### Installation Configuration Reference:

| Property                   | Description                                                         | Default           |
| -------------------------- | ------------------------------------------------------------------- | ----------------- |
| kafka.consumer.broker.host | Host or IP addresses of Kafka brokers for consumption of data       | none              |
| kafka.producer.broker.host | Host or IP addresses of Kafka brokers for publishing data           | none              |
| enable.checkpointing       | A boolean variable to enable checkpointing on cloud storage         | false             |
| redis.host                 | Host or IP address of Redis cache used for metadata caching         | none              |
| consumer.parallelism       | Number of threads to consume data in parallel                       | 1                 |
| operator.parallelism       | Number of threads to process data in parallel                       | 1                 |
| telemetry.schema.path      | Directory path for JSON schema files to validate the telemetry data | schemas/telemetry |
| event.max.size             | Acceptable size of each event in bytes                              | 1 Mb              |
| redis.devicestore.id       | The index of the device data store in Redis cache                   | 2                 |
| redis.userstore.id         | The index of the user data store in Redis cache                     | 12                |
| redis.contentstore.id      | The index of the content data store in Redis cache                  | 5                 |
| redis.dialcodestore.id     | The index of the dialcode data store in Redis cache                 | 6                 |

####

{% embed url="<https://github.com/project-sunbird/sunbird-data-pipeline>" %}
Data Pipeline source code
{% endembed %}


# Data Service

Data service comprises of a group of APIs that are used for configuring jobs and data exhausts - creation, updation, deletion etc. of jobs and exhausts are carried out via these APIs (CRUD for jobs)

#### Data Exhaust APIs:

Data exhaust apis provide the ability to create reports from pre-defined dataset configurations. The apis also provide the ability to create new datasets to generate reports.

![Data Exhaust/Report APIs](/files/y1gjiQGMFbu97CtUth95)

#### Key Features:

1. Create Datasets: Create new reusuable, custom datasets to generate custom reports. For instance, a dataset with an ability to apply a filter on course batchIds to generate reports.
2. Submit Reports: Ability to submit report requests based on pre-existing datasets. Currently supported datasets are:
   * Course Progress Reports
   * Userinfo Reports
   * Assessment Response Reports
3. Standard Datasets: Standard datasets do not have an ability to apply filters other than the tenant information and data ranges. Currently support standard datasets are:
   * Raw data reports
   * Summary data reports
   * Aggregated Summary data reports
4. Public Datasets: The public dataset api can be used to expose a scheduled report to public without authentication.

{% hint style="info" %}
[Data Exhaust API Documentation](http://docs.sunbird.org/latest/apis/dataexhaustapi/index.html)
{% endhint %}

#### Report APIs:

The Report APIs provide the ability to generate and schedule custom reports from Druid data store. The APIs leverage the query infrastructure provided by the Druid data store to configure custom reports

![](/files/y1S5sYYG99SlqX3N4RX2)

#### Key Features:

1. Schedule reports: The Report APIs provides the capability to execute a configured report on a schedule. The reports are uploaded to a configurable cloud storage.
2. Update pre-configured reports: The Report APIs provides the capability to update the report configuration of an existing scheduled report.
3. Configurable report configuration: The Reports can be set up with configurable queries to generate custom report. The report queries can have dimensions, aggregations, filters etc.

{% hint style="info" %}
[Report APIs Documentation](http://docs.sunbird.org/latest/apis/druidreportapi/index.html)
{% endhint %}

**Additional Documentation:**

**Accessing Sunbird Data Exhaust APIs :**

{% file src="/files/FrgkmH6tJmq7VrzOrC6W" %}

{% embed url="<https://github.com/project-sunbird/sunbird-analytics-service>" %}
Data Service source code
{% endembed %}

\\


# Data Product


# On Demand Druid Exhaust Job

On Demand Druid Exhaust Job is a generic data-product used to generate CSV reports. By passing the druid query config, we can use its capability to generate reports dynamically for any columns included in the druid datasource.


# Component Diagram

<figure><img src="/files/U0kE4DmGj7RcaKyeAyL2" alt=""><figcaption></figcaption></figure>

On Demand Druid Exhaust service will generate CSV reports based on user request. As this is a generic data-product user can request a CSV report for selected columns using filters.

1. **Database Layer**:

* **PostgreSQL Database (job\_request)**: This is where job requests are stored. These requests include information about job configuration. This data appended to postgress by [Filter Format From UI](https://project-sunbird.atlassian.net/wiki/spaces/MC/pages/3339780108/PD+-+Form+Config+-+Release-6.0.0+User+Detail+Report).
* **Druid:** This is where flattened data is stored with the help [ml-analytics](https://ed.sunbird.org/contribute/source-code/workflows/manage-learn/ml-analytics-service/ingestions) ingestion specs and can retrieve data using Druid queries for specific datasorces using [Model Config](https://github.com/Sunbird-Ed/ml-analytics-service/blob/release-5.1.0/druid_data_product_query_config.txt).\
  \
  **Data provider**

| Database   | Table/Datasouces                                                               |
| ---------- | ------------------------------------------------------------------------------ |
| PostgreSQL | job\_request                                                                   |
| Druid      | sl-project, sl-observation, sl-observation-status, sl-survey, ml-survey-status |

2. **Data Processing Layer**:\
   Apache Spark is used to perform transformations, sort columns, eliminate duplicates, and replace unknown values with null. This process enhances data quality, organizes data logically before storing to CSV.

## **User Interaction Diagram**

<figure><img src="/files/RIl8x9lhETsz1ygFmS3v" alt=""><figcaption></figcaption></figure>

This interaction diagram details the complete process of requesting and generating reports. The user can request a specific report through SunbirdEd from the program dashboard. Using exhaust APIs, this will map the request to SunbirdObsrv. OnDemondDruidExhaust data-product will be triggered by a scheduled cron task, which will query postgress and druid to get data and process it using Spark to transform data and generate the report. The user receives the same report once it has been created.


# ML CSV Reports

[Manage Learn](https://ed.sunbird.org/learn/product-and-developers-guide/manage-learn) (ML) is a vertical within the Sunbird project that focuses on implementing project, observation, and survey capabilities.

1. [Project](https://ed.sunbird.org/learn/product-and-developers-guide/manage-learn/what-is-a-project) : The improvement project empowers leaders to outline and monitor a series of tasks, guiding users toward achieving specific improvement objectives.
2. [observation](https://ed.sunbird.org/learn/product-and-developers-guide/manage-learn/what-is-observation) : Observations consist of questionnaires tailored for specific entities like school blocks or clusters.
3. [Surveys](https://ed.sunbird.org/learn/product-and-developers-guide/manage-learn/what-is-a-survey): Surveys on the Managed Learn are designed to gather valuable opinions and feedback from users without being associated with any specific entity.

## Different types of ML reports

<figure><img src="/files/oLBcd5M0MTZMooQSiOYi" alt=""><figcaption></figcaption></figure>

### **Project Reports**

1. **Task Detail Report :** This CSV report contains information related detailed a project progress along with task, sub-task data and evidences attached to tasks once completed. This report queries *sl-project* Druid datasource. This report is generated considering several [scenarios](https://docs.google.com/spreadsheets/d/1yE6G6sugfHiTvIWvl4AazShNr4kUNUcxqWQvZBpghsU/edit#gid=0).

* Sample **Task Detail Report** CSV

| UUID                                 | User Type     | User sub type | Declared State | District  | Block    | School Name              | School ID   | Declared Board     | Org Name                       | Program Name        | Program ID                     | Project ID               | Project Title                              | Project Objective                          | Category                           | Project start date of the user | Project completion date of the user | Project Duration | Project Status | Tasks    | Sub-Tasks             | Task Evidence                                                                                                                                                                             | Task Remarks | Project Evidence                                                                                                                                                                          | Project Remarks                   |
| ------------------------------------ | ------------- | ------------- | -------------- | --------- | -------- | ------------------------ | ----------- | ------------------ | ------------------------------ | ------------------- | ------------------------------ | ------------------------ | ------------------------------------------ | ------------------------------------------ | ---------------------------------- | ------------------------------ | ----------------------------------- | ---------------- | -------------- | -------- | --------------------- | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ------------ | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | --------------------------------- |
| 58b1a1a1-face-492b-a75a-0503863722bc | administrator | HM            | Andhra Pradesh | ANANTAPUR | AMADAGUR | MPUPS KANDUKURIVARIPALLI | 28224500205 | State (Tamil Nadu) | Staging Custodian Organization | Testing program 4.7 | PGM-FD255-testing\_program-4.7 | 6310918982642d0007cb71cb | ' Project with Tasks and Subtasks - FD255' | ' Project with Tasks and Subtasks - FD255' | 'Education Leader, School Process' | 2022-09-01T11:03:37.895+05:30  | 2022-09-01T11:06:53.582+05:30       | 2 weeks          | submitted      | 'Task 1' | 'subtask 1 of Task 1' | <https://sunbirdstagingpublic.blob.core.windows.net/samiksha/survey/6310918982642d0007cb71cb/58b1a1a1-face-492b-a75a-0503863722bc/293821ac-af9d-4853-89c2-1875a1f1fb9e/1662030251735.jpg> | Null         | <https://sunbirdstagingpublic.blob.core.windows.net/samiksha/survey/6310918982642d0007cb71cb/58b1a1a1-face-492b-a75a-0503863722bc/293821ac-af9d-4853-89c2-1875a1f1fb9e/1662030337167.jpg> | 'awesome content with resources ' |

2. **Status Report :** This CSV report contains information related to the status of the project progress user wise. This report queries sl-project Druid datasource.

* Sample **Status Report** CSV

| UUID                                 | User Type     | User sub type | state\_name    | District | Block       | School Name         | School ID   | Declared Board         | Org Name                       | Program Name        | Program ID                     | Project ID               | Project Title                              | Project Objective                          | Project start date of the user | Project completion date of the user | Project Duration | Last Synced date | Project Status | Certificate Status |
| ------------------------------------ | ------------- | ------------- | -------------- | -------- | ----------- | ------------------- | ----------- | ---------------------- | ------------------------------ | ------------------- | ------------------------------ | ------------------------ | ------------------------------------------ | ------------------------------------------ | ------------------------------ | ----------------------------------- | ---------------- | ---------------- | -------------- | ------------------ |
| 0b241dc8-c66e-4909-b4dd-ca1dbde28bd7 | administrator | HM            | Andhra Pradesh | CHITTOOR | B.KOTHAKOTA | MPPS GUTTAMEEDA H/W | 28230500607 | State (Andhra Pradesh) | Staging Custodian Organization | Testing program 4.7 | PGM-FD255-testing\_program-4.7 | 624168e6ed274f000880e853 | ' Project with Tasks and Subtasks - FD255' | ' Project with Tasks and Subtasks - FD255' | 2022-03-28T07:51:02.230+05:30  | Null                                | 2 weeks          | Null             | started        | Null               |

3. **Filtered Task Detail Report :** This CSV report contains information related detailed to project progress along with task, sub-task data and evidences attached to tasks once completed. Users can filter the volume of data by using the inbuild filters provided in the Program Dashboard. This report queries sl-project Druid datasource. This report is generated considering few [scenarios](https://docs.google.com/spreadsheets/d/1yE6G6sugfHiTvIWvl4AazShNr4kUNUcxqWQvZBpghsU/edit#gid=0).

* Sample **Filtered Task Detail Report** CSV

| UUID                                 | User Type     | User sub type | Declared State | District  | Block | Org Name                       | Program Name        | Project Title                              | Project Objective                          | Project Status | Project Completion Date of the user | Tasks    | Task Evidence                           | Task Remarks | Project Evidence                        | Project Remarks |
| ------------------------------------ | ------------- | ------------- | -------------- | --------- | ----- | ------------------------------ | ------------------- | ------------------------------------------ | ------------------------------------------ | -------------- | ----------------------------------- | -------- | --------------------------------------- | ------------ | --------------------------------------- | --------------- |
| 3486fdb5-0e13-4638-94af-9a0e5fd7d865 | administrator | DEO           | Andhra Pradesh | ANANTAPUR | AGALI | Staging Custodian Organization | Testing program 4.7 | ' Project with Tasks and Subtasks - FD255' | ' Project with Tasks and Subtasks - FD255' | submitted      | 2022-03-23T06:50:35.793+05:30       | 'Task 1' | [www.google.com](http://www.google.com) | Null         | [www.google.com](http://www.google.com) | "good work"     |

***

### **Survey Reports**

1. **Question Report :** This CSV report contains information of the questions attended by the users, their responses, remarks and evidence if submitted. This report queries sl-survey Druid datasource.

* Sample **Question Report** CSV

| UUID                                 | User Type     | User Sub Type | Declared State | District     | Block    | School ID   | School Name         | Declared Board | Organisation Name              | Program Name               | Program ID                  | Survey Name | Survey ID                                          | survey\_submission\_id   | Question\_external\_id                      | Question                  | Question\_response\_label | Evidences | Remarks                                                                                                                               |
| ------------------------------------ | ------------- | ------------- | -------------- | ------------ | -------- | ----------- | ------------------- | -------------- | ------------------------------ | -------------------------- | --------------------------- | ----------- | -------------------------------------------------- | ------------------------ | ------------------------------------------- | ------------------------- | ------------------------- | --------- | ------------------------------------------------------------------------------------------------------------------------------------- |
| 12dd9732-841c-47bc-8714-423e7e723961 | administrator | HM,DEO,SPD    | Andhra Pradesh | Ananthapuram | AMADAGUR | 28224500612 | MPPS(URDU) AMADAGUR | CBSE           | Staging Custodian Organization | Program - HT and officials | Program\_HT\_officials\_123 | Survey HT   | 54468f4e-14d8-11ee-af5f-bb2c8e1fa21b-1687862993715 | 64b67ec9c60fd1000869717c | SUR\_TEST\_001\_1687862979157-1687862993720 | Enter your First question | 'manual testing '         | Null      | 'test12333@#₹\_&-+()/?!;::''\*€¥$¢^°°={}\\\∆§×π√•\|\`\~\~. test. test ,,,, ...... test ,😣😀😕😕😕😕😕😭😭 test. test.....,,,,, test' |

2. **Status Report :** This CSV report contains the status information of Surveys attended by the users. This report queries ml-survey-status Druid datasource.

* Sample **Status Report** CSV

| UUID                                 | User type     | user\_sub\_type | Declared State | District                     | District  | School ID   | School Name                | Declared Board | Org Name                       | Program Name               | Program ID                  | Survey Name | Survey ID                                          | survey\_submission\_id   | Status of submission | Submission date |
| ------------------------------------ | ------------- | --------------- | -------------- | ---------------------------- | --------- | ----------- | -------------------------- | -------------- | ------------------------------ | -------------------------- | --------------------------- | ----------- | -------------------------------------------------- | ------------------------ | -------------------- | --------------- |
| 03f2e690-c031-4e68-aec7-521afae36bf0 | administrator | HM,DEO          | Andhra Pradesh | Dr. B.R. Ambedhkar Konaseema | Ainavilli | 28144900701 | MPPS (NO.2) SANAPALLILANKA | CBSE           | Staging Custodian Organization | Program - HT and officials | Program\_HT\_officials\_123 | Survey HT   | 54468f4e-14d8-11ee-af5f-bb2c8e1fa21b-1687862993715 | 649c1d9d5c81330008b7ea64 | started              | Null            |

***

### Observation with Rubric **Reports**

1. **Question Report :** This CSV report contains information of the questions attended by the users, their responses, remarks and evidence if submitted. This report queries sl-observation Druid datasource.

* Sample **Question Report** CSV

| UUID                                 | User Type | User Sub Type | Declared State | District | Block | School ID | School Name   | Declared Board | Org Name      | Program Name | Program ID             | Observation Name                 | Observation ID                                                   | Entity Observed | observation\_submission\_id | Domain Name | Criteria Name        | Question\_external\_id           | Question                                                                                                                                             | Question\_response\_label                            | Question score | Evidences | Remarks |
| ------------------------------------ | --------- | ------------- | -------------- | -------- | ----- | --------- | ------------- | -------------- | ------------- | ------------ | ---------------------- | -------------------------------- | ---------------------------------------------------------------- | --------------- | --------------------------- | ----------- | -------------------- | -------------------------------- | ---------------------------------------------------------------------------------------------------------------------------------------------------- | ---------------------------------------------------- | -------------- | --------- | ------- |
| 0498cee8-b015-4e51-a3a5-1bdc91c26914 | Null      | HM,DEO,SPD    | Andhra Pradesh | NELLORE  | KOVUR | Null      | ZPHS(G) KOVUR | CBSE           | ZPHS(G) KOVUR | Testing 4.4  | PGM\_FD\_98\_TEST\_4.4 | Observation with Rubrics – FD 98 | ce643538-332f-11ec-988d-cf2eb532a059-OBSERVATION-TEMPLATE\_CHILD | Andhra Pradesh  | 6295df3cae13230007dda6cb    | HP Domain   | Coaching & mentoring | Q17\_1634904093354-1634906556742 | You have to conduct regular review meetings with school principals of your district. How would you ensure that the review meeting is conducted well? | Make your own observations and ask general questions | 25             | Null      | Null    |

2. **Status Report :** This CSV report contains the status information of observation attended by the users. This report queries sl-observation-status Druid datasource.

* Sample **Status Report** CSV

| UUID                                 | User Type     | User sub type | Declared District | Declared District | Declared Block | Declared School ID | Declared School Name | Declared Board | Org Name                       | Program Name | Program ID             | Observation Name                 | Observation ID                                                   | District observed | Block observed | School observed | ID of school observed | Observation Submission ID | Status of submission | Submission date               | ECM marked NA |
| ------------------------------------ | ------------- | ------------- | ----------------- | ----------------- | -------------- | ------------------ | -------------------- | -------------- | ------------------------------ | ------------ | ---------------------- | -------------------------------- | ---------------------------------------------------------------- | ----------------- | -------------- | --------------- | --------------------- | ------------------------- | -------------------- | ----------------------------- | ------------- |
| 0498cee8-b015-4e51-a3a5-1bdc91c26914 | administrator | HM,DEO,SPD    | Andhra Pradesh    | NELLORE           | KOVUR          | 28192600523        | ZPHS(G) KOVUR        | CBSE           | Staging Custodian Organization | Testing 4.4  | PGM\_FD\_98\_TEST\_4.4 | Observation with Rubrics – FD 98 | ce643538-332f-11ec-988d-cf2eb532a059-OBSERVATION-TEMPLATE\_CHILD | Null              | Null           | Null            | Null                  | 6295df3cae13230007dda6cb  | completed            | 2022-05-31T09:29:03.491+05:30 | Null          |

3. **Domain Criteria Report :** This report contains questions and answer related information along with calculated the scores of domains. Each domain will contain criterias and each critiria contains numbers with score, so based on that it will calculate the scores. This report queries sl-observation-status sl-observation Druid datasource.

* Sample **Domain Criteria Report** CSV

| UUID                                 | User Type | User Sub Type | Declared State | District | Block | School ID | School Name   | Declared Board | Org Name      | Program Name | Program ID             | Observation Name                 | Observation ID                                                   | Entity Observed | observation\_submission\_id | Domain Name | Domain Level | Criteria Name | Criteria Level |
| ------------------------------------ | --------- | ------------- | -------------- | -------- | ----- | --------- | ------------- | -------------- | ------------- | ------------ | ---------------------- | -------------------------------- | ---------------------------------------------------------------- | --------------- | --------------------------- | ----------- | ------------ | ------------- | -------------- |
| 0498cee8-b015-4e51-a3a5-1bdc91c26914 | Null      | HM,DEO,SPD    | Andhra Pradesh | NELLORE  | KOVUR | Null      | ZPHS(G) KOVUR | CBSE           | ZPHS(G) KOVUR | Testing 4.4  | PGM\_FD\_98\_TEST\_4.4 | Observation with Rubrics – FD 98 | ce643538-332f-11ec-988d-cf2eb532a059-OBSERVATION-TEMPLATE\_CHILD | Andhra Pradesh  | 6295df3cae13230007dda6cb    | HP Domain   | L1           | Influence     | L3             |

***

### Observation without Rubric **Reports**

1. **Question Report :** This CSV report contains information of the questions attended by the users, their responses, remarks and evidence if submitted. This report queries sl-observation Druid datasource.

* Sample **Question Report** CSV

| UUID                                 | User Type | User Sub Type | Declared State | Declared District | Declared Block | Declared School ID | Declared School Name | Declared Board | Organisation Name              | Program Name      | Program ID         | Observation Name      | Observation ID                                                   | District observed | Block observed | School observed | ID of school observed | observation\_submission\_id | Submission date          | Question\_external\_id          | Question                      | Question\_response\_label | Question score | Evidences | Remarks |
| ------------------------------------ | --------- | ------------- | -------------- | ----------------- | -------------- | ------------------ | -------------------- | -------------- | ------------------------------ | ----------------- | ------------------ | --------------------- | ---------------------------------------------------------------- | ----------------- | -------------- | --------------- | --------------------- | --------------------------- | ------------------------ | ------------------------------- | ----------------------------- | ------------------------- | -------------- | --------- | ------- |
| 7e392af9-4923-4d14-8c3e-a2e8aab51237 | teacher   | TEACHER       | Andhra Pradesh | ANANTAPUR         | AMADAGUR       | 28224500518        | MPPS JOWKULA         | CBSE           | Staging Custodian Organization | Program Teacher-1 | Pgm\_Teacher\_2-QA | obs without rubrics 4 | be7b6ed0-0c35-11ee-9eca-334a9e787416-OBSERVATION-TEMPLATE\_CHILD | Anantapur         | AMADAGUR       | MPPS JOWKULA    | 28224500518           | 64993c9a5c81330008b78a49    | 2023-06-26T07:22:27.385Z | Q1\_1686913540822-1686913551805 | Enter the date of observation | 2023-06-26T12:52:00+05:30 | Null           | Null      | Null    |

2. **Status Report :** This CSV report contains the status information observation attended by the users. This report queries sl-observation-status Druid datasource.

* Sample **Status Report** CSV

| UUID                                 | User Type     | User sub type | Declared State | Declared District | Declared Block | Declared School ID | Declared School Name | Declared Board         | Org Name                       | Program Name      | Program ID         | Observation Name      | Observation ID                                                   | District observed | Block observed | School observed     | ID of school observed | Observation Submission ID | Status of submission | Submission date |
| ------------------------------------ | ------------- | ------------- | -------------- | ----------------- | -------------- | ------------------ | -------------------- | ---------------------- | ------------------------------ | ----------------- | ------------------ | --------------------- | ---------------------------------------------------------------- | ----------------- | -------------- | ------------------- | --------------------- | ------------------------- | -------------------- | --------------- |
| 52907735-7316-4d8d-886d-af46a8830549 | administrator | HM,DEO        | Andhra Pradesh | Ananthapuram      | Agali          | 28226200910        | MPPS HANUMANNAHALLI  | State (Andhra Pradesh) | Staging Custodian Organization | Program Teacher-1 | Pgm\_Teacher\_2-QA | obs without rubrics 4 | be7b6ed0-0c35-11ee-9eca-334a9e787416-OBSERVATION-TEMPLATE\_CHILD | Sri Satyasai      | Agali          | MPPS HANUMANNAHALLI | 28226200910           | 649abbb35c81330008b7c455  | started              | Null            |


# Folder Struture

On Demand Druid Exhaust Job Service folder structure is designed to organize the different modules and files. It follows a modular approach, facilitating easy management and development of the service.

### The structure is as follows

<pre><code>.
├── main    
<strong>│   └──  exhaust
</strong>│       ├── OnDemandBaseExhaustJob.scala
│       └── OnDemandDruidExhaustJob.scala
└── Test    
    └── exhaust
        └── TestOnDemandDruidExhaustJob.scala
</code></pre>

**main.exhaust**

This main directory houses our On Demand Druid Exhaust data-product. The *OnDemandDruidExhaustJob.scala* file contains main method which internally triggers other functions. *OnDemandBaseExhaustJob.scala* file contains generic functions such as execution, data transformation, storage in blob storage, and retrieval of file paths.

**test.exhaust**

In the exhaust folder we have the main test case file called *TestOnDemandDruidExhaustJob.scala* that is used to test our data-product.

### **Source code**

{% @github-files/github-code-block %}


# Report Service

The Report Service provides capabilities to create and consume reports on the front user interfaces/user apps. The reports are rendered and managed through report configurations.

![](/files/g7nSYrXad0Wjxj5Jy48f)

#### Key Features:

1. Scalable rendering: The report configuration and data files are rendered from a cloud storage or a CDN which supports large number of concurrent users accessing the reports from the front end user interface.
2. Easy to update: The update api allows the report configurations to be updated requiring no downtime for the reporting system.
3. Decoupling: Visualization (charts) and data are decoupled. This allows the system to reuse the same data for multiple visualizations.\
   \\

{% hint style="info" %}
[Report Service Documentation](http://docs.sunbird.org/latest/apis/reports/)
{% endhint %}

{% embed url="<https://github.com/project-sunbird/sunbird-report-service>" %}
Report Service source code
{% endembed %}


# Report Configurator

A self-service visualization tool for managing the lifecycle of reports (creation, review and publishing) and report configurations

![Report Configurator Workflow](/files/EslWs1eGGrRjYM91fmJJ)

#### Key Features:

1. Rich visualization support: Apache Superset is a open-sourced, rich, visualization tool to explore and visualize their data. Superset is supported out of the box in Sunbird Analytics to visualize and configure reports.
2. Access controls: Configurable access controls to create and publish reports to the backend report processor. The access for various functionalities can be controlled using user roles.
3. Configurable Storage: The generated reports can be configured on any cloud storage of choice and is flexible to support popular cloud storage such as AWS and Azure.

#### Additional Documentation:

Sunbird Report Configuration Guide:

{% file src="/files/MbNXYYb3oEr8vyVGG6c8" %}

**Source Code:**

{% embed url="<https://github.com/Sunbird-Ed/incubator-superset>" %}
Customized Report Publisher
{% endembed %}


# Summarisers

A generic summariser is generated from the telemetry events to enable easier access of frequently accessed metrics. For e.g., session level summaries or content player level summaries are some of the frequently accessed metrics derived for an actor/user. The content player level metrics can possibly be used to compute the total time spent by the user by playing various content in the platform across multiple sessions in a day.

```
app start
  session start
  ____
    workflow start (course start)
    ____
    ____
      play start
      ____
      ____
      play end
      ____
      play start
      ____
      ____
      play end
    ____
    ____
    workflow end
  ____
  session end
app end

Summary counts:
- App Summaries: 1
- Session Summaries: 1
- Workflow Summaries: 1
- Player Summaries: 2
```

#### Key Features:

1. Frequently used metrics: The generic summariser generates some of the frequently used metrics which can be used to generate custom data products for deriving various insights.
2. Data size reduction: Custom data products need to traverse huge amount of data to derive insights if they were to operate on the raw telemetry events generated from various devices. The summariser reduces the data traversal to almost 1/100th by summarising the data at a device level. This reduces the computation time and resources as well.

{% embed url="<https://github.com/project-sunbird/sunbird-core-dataproducts>" %}
Data Products source code
{% endembed %}


# ENGAGE




---

[Next Page](/llms-full.txt/1)

