1 of 48

DEVELOP AN OPTIMIZED FRAMEWORK FOR ANALYSIS OF STATIC AND RAPID IOT DATA

RESEARCH & DEVELOPMENT PROJECT

2 of 48

Recap

Introduction

Literature Review

Proposed Methodology

Tech Stack

Workflow

Results And Analysis

Challenges So Far

Future Scope

References

TABLE OF CONTENT

3 of 48

Completed Extract And First Phase of Transform

Analysis On Sample Data

RECAP

ETL Process

4 of 48

INTRODUCTION

5 of 48

INTRODUCTION

2.

CHALLENGE

Traditional data analysis techniques lead to delays and inaccuracies

3.

SOLUTION

Need framework to process massive data

4.

WHY

Analyze both static and real time data. To ensure accuracy

1.

IoT

By analyzing this data, significant insights can be gained.

6 of 48

Developing an optimized framework for analyzing static and rapidly changing IoT And Provide Real-Time Insights Into The Data Generated By IOT Devices

PROBLEM STATEMENT

7 of 48

LITERATURE REVIEW

8 of 48

RELATED

WORKS

9 of 48

GAPS

Papers reviewed recognize the challenges in processing real time data.

01

They rely on the traditional transformation steps and apply these to current ETL/ELT tools to deliver data to the front end, leading to failures.

02

They DO NOT explore the natural ability of stream data to be grouped under “windows” either during

  1. Extraction from multiple sources, or
  2. Transformation, or
  3. Loading

03

10 of 48

PROPOSED METHODOLOGY

11 of 48

OVERVIEW

Loading the final data into Hive Databases

05

Defining the Problem and Objectives

01

Data Transformation

04

Data Collection and Preprocessing

02

Windowing

03

12 of 48

DEFINING THE PROBLEM AND OBJECTIVES

Data Collection & Preprocessing

Data pipeline to collect real-time streaming data from IoT devices

Using Morphline

Flume configuration files for data sources, channels, and sinks

02

-

Helps with ambiguous data

System designed to effectively gather and transfer streaming data.

Defining Problem & Objectives

To retrieve data from IoT devices, we employed Apache Flume

01

-

13 of 48

WINDOWING

Windowing based on local timestamps

Dividing the data into smaller time intervals or windows for better querying and transformation

01

02

Using conf files

Utilizing Morphline properties to pipeline the data.

morphlines.conf

morphex.conf

morph_result.conf

14 of 48

DATA TRANSFORMATION

Partitioning and Clustering

Partitioning - Divides data into directories for easier data management.

Clustering - Organises data based on columns to improve query performance.

Using HIVE

Windowed data to be transformed into a Unified Schema using HIVE’s properties.

01

02

15 of 48

LOADING IN DATABASES

Ease of efficient querying

It enables the integration of data from multiple sources for efficient querying and analysis.

Stores Data on HDFS

Unified Schema

Unified Schema is a standardized way to store and access data in HIVE databases.

01

02

Creation of HIVE Database

Starting by creating a HIVE table using Command Line Interface

03

16 of 48

17 of 48

TECH STACK

18 of 48

TECH STACK

HDFS

Map-Reduce

01

HADOOP

Agents

Data-Flow

02

FLUME

Metastore

Query Engine

03

HIVE

19 of 48

WORKFLOW

20 of 48

INPUT SANITIZATION

The following structure of data gives parsing error due to Grok Expectations

Hence the input was correctly sanitized to be parsed by the upcoming stages

21 of 48

STAGE 1: DATA INPUT AND PARSING

MORPH_RESULT.CONF

22 of 48

STAGE 2: DATA INPUT AND PARSING

MORPHEX.CONF

23 of 48

LOADING OF DATA IN HIVE

timewindow

stage2

stage1

data

24 of 48

LOADING OF DATA IN HIVE

time_05_05_2017/stream.1683606675748

data

Processed Data

s1flumeDir1

25 of 48

EXTRACTION IN HIVE

stage1

data

Split data

26 of 48

TRANSFORMATION OF DATA IN HIVE

Timestamped data

stage2

stage1

27 of 48

PARTITIONING OF DATA IN HIVE

Final data

timewindow

stage2

28 of 48

FINAL HIVE TABLES IN HDFS

user/hive/warehouse/stage1

user/hive/warehouse/stage2

user/hive/warehouse/timewindow

29 of 48

RESULTS AND ANALYSIS

30 of 48

THE CONFIGURATION FILES

Structure for the Stage 1

Morph_result.conf file is used in our Morphline use-case to define a proper structure for the stage 1 of the development

MORPH_RESULT.CONF

Utilizes morphline interceptors

This file processes the input data and puts it into the results folder for the next stage.

02

01

31 of 48

THE CONFIGURATION FILES

Structure for FLUME agent for stage 2

Morphex.conf file is used in our Morphline use-case to define a proper structure for our Flume Agent to work.

MORPHEX.CONF

Defining agent and it’s working

It defines the Source, the Sink, and the channels and their types to facilitate data transfer between both of them

02

01

32 of 48

THE CONFIGURATION FILES

Definition

Morphline.conf is our configuration file used in Flume to define ETL transformations on data streams.

MORPHLINE.CONF

Parsing of Data

It specifies a series of commands to process and modify data coming in from the source.

02

01

33 of 48

HADOOP FILE SYSTEM

All the data is first processed through our 2 phases and then stored in the HDFS directory s1flumeDir1

After the preprocessing is complete the data is arranged in the following tables through HIVE, making the framework more scalable

34 of 48

ANALYSIS

SCALABLE

Partitioning

Clustering

01

FLEXIBLE

Morphex.conf

Morphline.conf

02

Morph_result.conf

HIVE

FLUME

35 of 48

ANALYSIS

EFFICIENT

CUSTOM MONTH INTERCEPTOR

CUSTOM TIME INTERCEPTOR

36 of 48

CHALLENGES SO FAR

37 of 48

ISSUE : PARALLEL LAUNCH OF HIVE METASTORE NOT SUPPORTED

Resolution

Running the shown command separately resolved the issue

03

Parallel launch of HIVE Metastore not supported

01

38 of 48

HADOOP SAFE MODE ON

RESOLUTION

01

02

Hadoop "safemode" is essentially a restricting setting for the HDFS

Led the HIVE to not get started as it requires full data transaction

Manually fix this issue by separately forcing Hadoop as an admin to leave the safemode

39 of 48

ISSUE: UNRESOLVED IMPORTS

We got this issue of unresolved imports in the JAR files

01

02

Didn’t let the codes to execute properly

40 of 48

RESOLUTION: UNRESOLVED IMPORTS

Resolved by importing the required libraries

01

02

Libraries correctly added to the classpath of the newly exported JAR files

41 of 48

ISSUE: RECURRING WARNINGS

Infinite loop of the above WARNING messages.

01

02

Warn messages were occuring due to expectations of grok patterns.

42 of 48

RESOLUTION: INPUT SANITIZATION

01

02

We manually fixed the csv file by sanitising the input files as per the Grok.

By changing the structure of the data.

43 of 48

FUTURE PROSPECTS

44 of 48

SHORTCOMINGS

SCOPE OF IMPROVEMENT

Seamless Flume-To-Hive Integration For Automation

Data Visualization And Dashboards

Scaling Data Analysis For Variety of Datasets

Continuous Optimisation And Automation

Integration Of Machine Learning Algorithms

45 of 48

REFERENCES

46 of 48

REFERENCES

  1. Pareek, A., Zhang, B., & Khaladkar, B. (2019, June). A Demonstration of Striim A Streaming Integration and Intelligence Platform. In Proceedings of the 13th ACM International Conference on Distributed and Event-based Systems (pp. 236-239).

3. Adilah Sabtu, Nurulhuda Firdaus Mohd Azmi , Nilam Nur Amir Sjarif , et al. (2017, November 30). THE CHALLENGES OF EXTRACT, TRANSFORM AND LOAD (ETL) FOR DATA INTEGRATION IN NEAR REAL TIME ENVIRONMENT. Journal of Theoretical and Applied Information Technology, 95(22), 6314–6322.

4. Pravin Gokhe (2016, March 26). ENTERPRISE REAL-TIME INTEGRATION (1st ed., Vol. 1) [English].(pp 51-57)

47 of 48

REFERENCES

5. Tönjes, R., Barnaghi, P., Ali, M., Mileo, A., Hauswirth, M., Ganz, F., & Puiu, D. (2014, June). Real time IoT stream processing and large-scale data analytics for smart city applications. (pp. 24-28, 39-44)

6. Hoffman, S. (2015). Apache flume: Distributed log collection for hadoop. Birmingham Packt Publishing.(pp. 62-84)

8. Apache (2014) Apache Hadoop. http://hadoop.apache.org/.

9. Capriolo E, Wampler D, Rutherglen J (2012 December). Programming Hive. O’Reilly Media, Inc. (pp. 49-97)

7. Waas, F., Wrembel, R., Freudenreich, T., Thiele, M., Koncilia, C., & Furtado, P. (2013, April 1). On-Demand ELT Architecture for Right-Time BI. International Journal of Data Warehousing and Mining, 9(2), 21–38.

48 of 48

THANK YOU

Now we will move to the Demonstration