Using Oracle Stream Analytics
Using Oracle Stream Analytics
F18429-20
April 2025
Oracle GoldenGate Stream Analytics Documentation,
F18429-20
This software and related documentation are provided under a license agreement containing restrictions on use and
disclosure and are protected by intellectual property laws. Except as expressly permitted in your license agreement or
allowed by law, you may not use, copy, reproduce, translate, broadcast, modify, license, transmit, distribute, exhibit,
perform, publish, or display any part, in any form, or by any means. Reverse engineering, disassembly, or decompilation
of this software, unless required by law for interoperability, is prohibited.
The information contained herein is subject to change without notice and is not warranted to be error-free. If you find
any errors, please report them to us in writing.
If this is software, software documentation, data (as defined in the Federal Acquisition Regulation), or related
documentation that is delivered to the U.S. Government or anyone licensing it on behalf of the U.S. Government, then
the following notice is applicable:
U.S. GOVERNMENT END USERS: Oracle programs (including any operating system, integrated software, any
programs embedded, installed, or activated on delivered hardware, and modifications of such programs) and Oracle
computer documentation or other Oracle data delivered to or accessed by U.S. Government end users are "commercial
computer software," "commercial computer software documentation," or "limited rights data" pursuant to the applicable
Federal Acquisition Regulation and agency-specific supplemental regulations. As such, the use, reproduction,
duplication, release, display, disclosure, modification, preparation of derivative works, and/or adaptation of i) Oracle
programs (including any operating system, integrated software, any programs embedded, installed, or activated on
delivered hardware, and modifications of such programs), ii) Oracle computer documentation and/or iii) other Oracle
data, is subject to the rights and limitations specified in the license contained in the applicable contract. The terms
governing the U.S. Government's use of Oracle cloud services are defined by the applicable contract for such services.
No other rights are granted to the U.S. Government.
This software or hardware is developed for general use in a variety of information management applications. It is not
developed or intended for use in any inherently dangerous applications, including applications that may create a risk of
personal injury. If you use this software or hardware in dangerous applications, then you shall be responsible to take all
appropriate fail-safe, backup, redundancy, and other measures to ensure its safe use. Oracle Corporation and its
affiliates disclaim any liability for any damages caused by use of this software or hardware in dangerous applications.
Oracle®, Java, MySQL, and NetSuite are registered trademarks of Oracle and/or its affiliates. Other names may be
trademarks of their respective owners.
Intel and Intel Inside are trademarks or registered trademarks of Intel Corporation. All SPARC trademarks are used
under license and are trademarks or registered trademarks of SPARC International, Inc. AMD, Epyc, and the AMD logo
are trademarks or registered trademarks of Advanced Micro Devices. UNIX is a registered trademark of The Open
Group.
This software or hardware and documentation may provide access to or information about content, products, and
services from third parties. Oracle Corporation and its affiliates are not responsible for and expressly disclaim all
warranties of any kind with respect to third-party content, products, and services unless otherwise set forth in an
applicable agreement between you and Oracle. Oracle Corporation and its affiliates will not be responsible for any loss,
costs, or damages incurred due to your access to or use of third-party content, products, or services, except as set forth
in an applicable agreement between you and Oracle.
Contents
1 Overview
1.1 Introduction 1-1
1.2 Key Features of GGSA 1-1
1.3 GGSA Architecture 1-2
1.4 Steps to build Continuous-ETL and Realtime-Analytics Pipelines 1-3
2 Install
2.1 Planning Your Installation 2-1
2.2 Installing GoldenGate Stream Analytics 2-3
2.3 Configuring the Metadata Store 2-4
2.3.1 Configuring ATP/ADW as Metadata Store 2-7
2.4 Initializing Metadata Store 2-8
2.5 Jetty Properties File 2-10
2.6 Adjusting Jetty Threadpool 2-11
2.7 Integrating Stream Analytics with Oracle GoldenGate 2-11
2.8 Maven Setting for GoldenGate Big Data Handlers 2-12
2.8.1 Set the Maven Home Path 2-12
2.8.2 Configure Maven Proxy Settings 2-12
2.9 GoldenGate Stream Analytics Hardware Requirements for Enterprise Deployment 2-13
2.10 Retaining https and Disabling http 2-16
2.11 Setting up Runtime for GoldenGate Stream Analytics Server 2-16
2.12 Validating Data Flow to GoldenGate Stream Analytics 2-19
2.13 Terminating GoldenGate Stream Analytics 2-20
2.14 Upgrading GoldenGate Stream Analytics 2-20
3 Configure
3.1 Configure Runtime Environment 3-1
3.1.1 Mandatory Configurations 3-1
[Link] Configuring Kafka 3-1
[Link] Configuring the Runtime Server 3-3
3.1.2 Optional Configurations 3-12
[Link] Configuring Pipeline Preferences 3-12
iii
[Link] Configuring Network Proxy 3-13
[Link] Configuring Kafka Preferences 3-14
[Link] Configuring GG Preferences 3-14
[Link] Configuring SQL Preferences 3-14
[Link] Changing Spark Work Directory 3-15
[Link] Changing Spark Log Rollover based on Time 3-15
3.2 Configure Users 3-16
3.2.1 Managing Users 3-16
[Link] Adding Users 3-17
[Link] Changing Password 3-19
[Link] Removing Users 3-19
[Link] Configuring LDAP for User Authentication and Management 3-20
3.2.2 Configuring User Preferences 3-22
4 Manage
4.1 Connections 4-1
4.1.1 Create Connections 4-1
[Link] Creating a Connection to ADW or ATP 4-2
[Link] Creating a Connection to AWS S3 4-3
[Link] Creating a Connection to Coherence 4-3
[Link] Creating a Connection to Druid 4-4
[Link] Creating a Connection to Elasticsearch 4-4
[Link] Creating a Connection to GoldenGate 4-5
[Link] Creating a Connection to HBase 4-6
[Link] Creating a Connection to HDFS 4-7
[Link] Creating a Connection to Hive 4-7
[Link] Creating a connection to Ignite Cache 4-8
[Link] Creating a Connection to JMS 4-9
[Link] Creating a Connection to Kafka 4-9
[Link] Creating a Connection to Microsoft Azure Data Lake-Gen2 4-10
[Link] Creating a Connection to MongoDB 4-11
[Link] Creating a Connection to MySQL Database 4-13
[Link] Creating a Connection to OCI Object Store 4-13
[Link] Creating a Connection to ONS 4-14
[Link] Creating a Connection to Oracle AQ 4-15
[Link] Creating a Connection to Oracle Database 4-16
[Link] Creating a Connection to OSS 4-16
4.1.2 Manage Connections 4-17
4.2 Streams 4-18
4.2.1 Create Streams 4-18
[Link] Creating a File Stream 4-18
iv
[Link] Creating a GoldenGate Stream 4-19
[Link] Creating a JMS Stream 4-21
[Link] Creating a Kafka Stream 4-23
4.2.2 Manage Streams 4-24
[Link] Application Timestamp 4-25
[Link] Supported Timestamp Formats in an Input Stream 4-25
[Link] Predefined CSV Data Formats 4-26
4.3 References 4-27
4.3.1 Create References 4-27
[Link] Creating a Coherence Reference 4-27
[Link] Creating a Database Reference 4-28
[Link] Creating an Ignite Reference 4-29
4.3.2 Manage References 4-30
[Link] Coherence Reference 4-31
4.4 Targets 4-34
4.4.1 Create Targets 4-34
[Link] Creating an AWS S3 Target 4-35
[Link] Creating an Azure DataLake Gen-2 Target 4-36
[Link] Creating a Coherence Target 4-37
[Link] Creating a Database Target 4-38
[Link] Creating an Elasticsearch Target 4-39
[Link] Creating an HBase Target 4-40
[Link] Creating HDFS Target 4-41
[Link] Creating a Hive Target 4-42
[Link] Creating an Ignite Cache Target 4-43
[Link] Creating a JMS Target 4-44
[Link] Creating a Kafka Target 4-46
[Link] Creating a MongoDB Target 4-47
[Link] Creating a Network File System (NFS) Target 4-48
[Link] Creating a Notification Target 4-49
[Link] Creating an OCI Object Store Target 4-50
[Link] Creating an OSS Target 4-51
[Link] Creating a REST Target 4-53
4.4.2 Manage Targets 4-55
[Link] Coherence Target 4-55
4.5 Pipelines 4-56
4.5.1 Create a Pipeline 4-57
4.5.2 Manage Pipelines 4-57
[Link] Using the Pipeline Editor 4-57
[Link] Publishing a Pipeline 4-57
[Link] Unpublishing a Pipeline 4-58
[Link] Exporting and Importing a Pipeline and Its Dependent Artifacts 4-58
v
[Link] Working with Live Output Table 4-60
[Link] Using the Topology Viewer 4-60
4.6 GoldenGate Change Stream 4-62
4.6.1 Getting a GoldenGate Change Stream into a Kafka Topic 4-62
4.6.2 Manage GG Change Data Stream 4-63
[Link] Starting a GoldenGate Change Stream 4-63
[Link] Stopping a GG Change Data Stream 4-64
[Link] Purging the GoldenGate Trail Files 4-64
[Link] Streaming GoldenGate Full Records 4-65
4.7 Embedded Ignite Cache 4-65
4.7.1 Starting a Cache Cluster 4-65
4.7.2 Stopping a Cache Cluster 4-66
4.7.3 Restarting a Cache Cluster 4-66
4.7.4 Monitoring Cache in the Cache Cluster 4-66
4.8 Ignite Cluster on OCI GGSA 4-67
4.8.1 Starting an Ignite Cluster 4-67
4.8.2 Scaling an Ignite Cluster 4-67
4.8.3 Deleting Storage 4-67
4.8.4 Stopping an Ignite Cluster 4-68
4.9 GGBD Cluster on OCI GGSA 4-68
4.9.1 Starting a GGBD Cluster 4-68
4.9.2 Stopping a GGBD Cluster 4-68
5 Transform
5.2 Correlating Streams and References 5-1
5.2.1 Joining Mutiple Streams 5-1
5.2.2 Joining a Stream with a Reference or an External Source 5-2
5.3 Applying Window Functions to a Stream 5-2
5.3.1 Applying a Time Window with Slide 5-2
5.3.2 Applying a Time Window without Slide 5-3
5.3.3 Applying a Row Window with Slide 5-3
5.3.4 Applying a Row Window without Slide 5-4
5.3.5 Applying a window with current year, month, day, or hour 5-4
5.3.6 Applying your own Window using Field from Payload 5-5
5.3.7 Applying a Row window with Partition without Range 5-5
5.3.8 Applying a Row Window with Partition with Range without Slide 5-5
5.3.9 Applying a Row Window with Partition with Slide and Range 5-6
5.1 Adding Stages to a Pipeline 5-6
5.1.1 Adding a Query Stage 5-6
5.1.2 Adding a Filter to a Query Stage 5-6
5.1.3 Adding a Summary to a Query Stage 5-7
vi
5.1.4 Adding a Summary with Group By 5-7
5.1.5 Adding a Query Group Stage 5-8
[Link] Adding Query Group: Stream 5-8
[Link] Adding Query Group: Table 5-9
5.1.6 Adding a Rule Stage 5-9
5.1.7 Adding a Pattern Stage 5-10
5.1.8 Adding a Scoring Stage 5-10
5.1.9 Adding a Target Stage 5-10
5.1.10 Adding a Custom CQL Stage 5-11
5.4 Applying Functions to Create a New Column 5-11
5.4.1 Using Bessel Functions 5-12
[Link] BesselI0 5-13
[Link] BesselIO_exp 5-13
[Link] BesselI1(value1) 5-13
[Link] BesselI1_exp(value1) 5-13
[Link] BesselK0_exp(value1) 5-14
[Link] BesselIK1_exp(value1) 5-14
[Link] BesselY(value1, value2) 5-14
[Link] BesselJ(value1, value2) 5-14
[Link] BesselK(value1,value2) 5-15
5.4.2 Using Conversion Functions 5-15
[Link] bigdecimal(value1) 5-15
[Link] boolean(value1) 5-15
[Link] double(value1) 5-16
[Link] float(value1) 5-16
[Link] int(value1) 5-16
[Link] long() 5-16
[Link] string(value1, value2) 5-17
5.4.3 Using Date Functions 5-17
[Link] Acceptable Formats for Timestamp Values 5-17
[Link] Day(date) 5-19
[Link] eventtimestamp(value1) 5-19
[Link] hour(date) 5-19
[Link] minute(date) 5-19
[Link] month(date) 5-20
[Link] nanosecond(value1) 5-20
[Link] systemtimestamp(value1) 5-20
[Link] timeformat(value1, value2) 5-20
[Link] Year(date) 5-21
5.4.4 Using Geometry Functions 5-21
[Link] CreatePoint(value1, value2, value3) 5-21
[Link] distance(lat1, long1, lat2, long2,SRID) 5-22
vii
5.4.5 Using Interval Functions 5-22
[Link] dsintervaltonum(value1, value 2) 5-23
[Link] numtodsinterval(value1, value2) 5-23
[Link] numtoyminterval(value1, value 2) 5-24
[Link] to_dsinterval(value1) 5-24
[Link] to_yminterval(value1) 5-24
[Link] ymintervaltonum(value1, value2) 5-25
5.4.6 Using Math Functions 5-25
[Link] IEEEremainder(value1, value1) 5-27
[Link] abs(value1) 5-27
[Link] acos(value1) 5-27
[Link] asin(value1) 5-27
[Link] atan(value1) 5-28
[Link] atan2 5-28
[Link] binomial(base, power) 5-28
[Link] bitMaskWithBitsSetFromTo(value1, value2) 5-28
[Link] cbrt() 5-29
[Link] ceil() 5-29
[Link] copySign() 5-29
[Link] cos(value1) 5-29
[Link] cosh(value1) 5-29
[Link] exp(value1, value2) 5-30
[Link] expm1(value1) 5-30
[Link] factorial(value1) 5-30
[Link] floor(value1) 5-30
[Link] GetExponent(value1) 5-30
[Link] getSeedAtRowColumn(value1, value2) 5-31
[Link] hash(value1) 5-31
[Link] hypot(value1, value2) 5-31
[Link] LeastSignificantBit(value1) 5-31
[Link] log(value1, value2) 5-32
[Link] log1(value1) 5-32
[Link] log10(value1) 5-32
[Link] log2(value1) 5-32
[Link] logFactorial(value1) 5-33
[Link] long() 5-33
[Link] longFactorial(value1) 5-33
[Link] minimum(value1, value2) 5-33
[Link] mod(value1, value2) 5-34
[Link] mostSignificantBit(value1) 5-34
[Link] nextAfter(value1, value2) 5-34
[Link] nextDown(value1, value2) 5-34
viii
[Link] nextUp(value1) 5-35
[Link] pow(value1, value2) 5-35
[Link] rint(value1) 5-35
[Link] round(value1) 5-35
[Link] scalb( 5-36
[Link] signum(value1) 5-36
[Link] sin(value1) 5-36
[Link] sinh(value1) 5-36
[Link] sqrt(value1) 5-37
[Link] stirlingCorrection(value1) 5-37
[Link] tan(value1) 5-37
[Link] tanh(value1) 5-37
[Link] toDegrees(value1) 5-37
[Link] toRadians(value1) 5-38
[Link] ulp(value1) 5-38
5.4.7 Using Null-related Functions 5-38
[Link] nvl(value1, value2) 5-38
5.4.8 Using Statistical Functions 5-39
[Link] beta1(value1, value2, value3) 5-40
[Link] betacomplemented(value1, value2, value3) 5-40
[Link] binomial2(value1, value2, value3) 5-40
[Link] binomialcomplemented(value1, value2, value3) 5-41
[Link] chiSquare(value1, value2) 5-41
[Link] chiSquareComplemented(value1, value2) 5-41
[Link] errorFunction(value1) 5-41
[Link] errorFunctionComplemented(value1) 5-42
[Link] gamma(value1, value2, value3) 5-42
[Link] gammacomplemented(value1, value2, value3) 5-42
[Link] incompleteBeta(value1, value2, value3) 5-43
[Link] incompleteGamma(value1, value2) 5-43
[Link] incompleteGammaComplement(value1, value2) 5-43
[Link] logGamma(value1) 5-43
[Link] negativeBinomial(value1, value2, value3) 5-44
[Link] negativeBinomialComplemented(value1, value2, value3) 5-44
[Link] normal(value1, value2, value3) 5-44
[Link] normalInverse(value1) 5-45
[Link] poisson(value1, value2) 5-45
[Link] poissonComplemented(value1, value2) 5-45
[Link] studentT(value1, value2) 5-45
[Link] studentTInverse(value1, value2) 5-46
5.4.9 Using String Functions 5-46
[Link] coalesce(value1,... ) 5-47
ix
[Link] Concat(value1,...) 5-47
[Link] indexof(value1, value2) 5-47
[Link] initcap(value1) 5-48
[Link] length(value1) 5-48
[Link] like(string, pattern) 5-48
[Link] lower(value1) 5-49
[Link] lpad(value1, value2, value3) 5-49
[Link] ltrim(value1, value2) 5-49
[Link] replace(string, match, replacement) 5-50
[Link] rpad(value1, value2, value3) 5-50
[Link] rtrim(value1, value2) 5-50
[Link] substr() 5-51
[Link] substring(string, from, to) 5-51
[Link] translate(expression, from_string, to_string) 5-51
[Link] upper(value1) 5-52
5.5 Adding Custom Functions and Custom Stages 5-52
5.5.1 Creating a Custom Jar 5-52
5.5.2 Adding Custom Functions 5-52
5.5.3 Implementing Custom Functions 5-53
[Link] Sample: Encrypt a Column 5-53
5.5.4 Adding a Custom Stage 5-53
[Link] Sample: Encrypt a Column 5-54
[Link] Sample: Invoke a REST Service 5-55
[Link] Sample: Invoke a SOAP Service 5-58
5.5.5 Limitations 5-60
5.5.6 Mapping of Data Types 5-61
5.6 Writing CQL Queries 5-61
5.6.1 Sample Queries 5-61
[Link] A Followed By B 5-62
[Link] A Not Followed by B 5-64
[Link] Detect Duplicates 5-64
[Link] Change Event 5-65
[Link] Eliminate Duplicates 5-66
6 Analyze
6.1 Using Geofences for Location-based Analytics 6-1
6.1.1 Selecting a Tile Layer 6-1
[Link] Elocation Tile Layer 6-1
[Link] Open Street Maps Tile Layer 6-2
[Link] Google Maps Tile Layer 6-3
[Link] Custom Tile Layer 6-4
x
6.1.2 Managing Geofences using the Map Editor 6-6
[Link] Creating a Geo Fence 6-6
[Link] Deleting a Geofence 6-7
6.1.3 Importing a Geofence from a Database 6-7
6.1.4 Using Spatial Patterns in Pipeline Stages 6-7
[Link] Clearing Objects Outside a Geo Fence 6-7
[Link] Tracking Objects using a Geo Fence 6-8
[Link] Getting Direction of a Moving Object 6-8
[Link] Obtaining Geographic Coordinates 6-9
[Link] Calculating Distance between Objects in a Stream 6-9
[Link] Calculating Distance between Objects in Two Streams 6-10
[Link] Creating Geo Fence 6-10
[Link] Monitoring Proximity between Objects in a Stream 6-10
[Link] Monitoring Proximity between Objects in Two Streams 6-11
[Link] Obtaining the Proximity of an Object from a Geo Fence 6-11
[Link] Finding Nearest Place using the Geographical Coordinates 6-12
[Link] Finding Nearest Place Details using the Geographical Coordinates 6-12
[Link] Determining Average Speed 6-13
6.2 Transforming and Analyzing Data using Patterns 6-13
6.2.1 Adding a Pattern Stage 6-15
6.2.2 Detecting Missing Events 6-15
6.2.3 Calculating Quantile Value 6-15
6.2.4 Identifying Correlation between Two Numeric Patterns 6-16
6.2.5 Detecting Duplicate Events 6-16
6.2.6 Eliminating Duplicate Events 6-17
6.2.7 Detecting Event Value Changes 6-17
6.2.8 Detecting Data Field Value Changes 6-18
6.2.9 Monitoring Sequence of Events 6-19
6.2.10 Outputting Highest Value Events 6-19
6.2.11 Outputting Lowest Value Events 6-20
6.2.12 Monitoring Invariably Increasing Numeric Values 6-20
6.2.13 Monitoring Invariably Decreasing Numeric Values 6-21
6.2.14 Identifying the Missing First Event in a Sequence 6-22
6.2.15 Identifying the Second Missing Event in a Sequence 6-22
6.2.16 Analyzing Data using Double Bottom Charts 6-23
6.2.17 Analyzing Data using Double Top Charts 6-23
6.2.18 Correlating Current and Previous Events 6-24
6.2.19 Delaying Delivery of Events to Downstream Node 6-25
6.2.20 Outputting Contents to Downstream Node 6-25
6.2.21 Outputting Unexpired Contents to Downstream Node 6-25
6.2.22 Merging Two Streams having Identical Shapes 6-26
6.2.23 Joining Flows with Streams and References 6-26
xi
6.2.24 Transforming Events into JSON 6-26
6.2.25 Transforming a Single Event from a Stage into Multiple Events 6-27
6.2.26 Merging Two Continuous Events into a Single Event 6-27
6.2.27 Applying OML Models to get the Scoring of Events (Preview Feature) 6-27
6.2.28 Detecting Contiguous Events 6-28
6.2.29 Creating Pivot Columns 6-28
6.3 Using Machine Learning Models for Scoring and Prediction 6-29
6.3.1 Importing a Predictive Model 6-29
6.3.2 Adding a Scoring Stage 6-29
6.4 Integrating with Druid Timeseries Database for Realtime Interactive Analytics 6-30
6.4.1 Creating a Connection to Druid 6-30
6.4.2 Creating a Cube 6-30
6.4.3 Exploring a Cube 6-32
7 Visualize
7.1 Adding Realtime Charts 7-1
7.1.1 Adding an Area Chart 7-1
7.1.2 Adding a Bar Chart 7-2
7.1.3 Adding a Bubble Chart 7-2
7.1.4 Adding a Line Chart 7-3
7.1.5 Adding a Pie Chart 7-4
7.1.6 Adding a Scatter Plot 7-4
7.1.7 Adding a Stacked Bar Chart 7-5
7.1.8 Adding a Thematic Map 7-5
7.1.9 Updating Visualizations 7-6
7.2 Creating and Managing Dashboards 7-6
7.2.1 Adding a Dashboard 7-6
7.2.2 Editing a Dashboard 7-7
7.2.3 Sharing a Dashboards with Peers 7-10
7.2.4 Deleting a Dashboard 7-10
7.2.5 Importing a Dashboard with all its Dependencies 7-10
7.2.6 Exporting a Dashboard with all its Dependencies 7-10
8 Monitor
8.1 Execution and HA Statistics 8-1
8.2 Detailed Query Analysis 8-3
8.3 Complete CQL Engine Statistics 8-4
xii
9 Reference
9.1 Pipeline Details 9-1
9.2 Stage Details 9-2
9.3 Query Details 9-3
9.4 Internal Kafka Topics 9-4
10 Troubleshoot
10.1 Pipeline Debug and Monitoring Metrics 10-1
10.1.1 Spark Standalone 10-1
10.1.2 Spark on YARN 10-1
10.1.3 Pipeline Details 10-3
10.1.4 Stage Details 10-4
10.1.5 Query Details 10-5
10.1.6 Execution and HA Statistics 10-6
10.1.7 Detailed Query Analysis 10-8
10.1.8 Complete CQL Engine Statistics 10-9
10.1.9 Internal Kafka Topics 10-10
10.2 Common Issues and Remedies 10-10
10.2.1 Pipeline 10-11
10.2.2 Pipeline 10-11
[Link] Pipelines are not running as expected 10-11
[Link] GGSA Pipeline getting Terminated 10-12
[Link] Live Table Shows Listening Events with No Events in the Table 10-12
[Link] Live Table Still Shows Starting Pipeline 10-13
[Link] Time-out Exception in the Spark Logs when you Unpublish a Pipeline 10-14
[Link] Piling up of Queued Batches in HA mode 10-14
[Link] Null Record from Summary in Query Stage 10-14
10.2.3 Stream 10-15
[Link] Cannot See Any Kafka Topic or a Specific Topic in the List of Topics 10-15
[Link] Input Kafka Topic is Sending Data but No Events Seen in Live Table 10-15
10.2.4 Connection 10-15
[Link] Database Connection Failure 10-15
[Link] Druid Connection Failure 10-16
[Link] Coherence Connection Failure 10-16
[Link] JNDI Connection Failure 10-16
10.2.5 Target 10-16
[Link] Cannot see any Events in Targets 10-17
10.2.6 Geofence 10-17
[Link] Name and Description Fields are not displayed for the DB-based
Geofences 10-17
[Link] DB-based Geofence is not Working 10-17
xiii
10.2.7 Cube 10-17
[Link] Unable to Explore Cube which was Working Earlier 10-17
[Link] Cube Displays "Datasource not Ready" 10-17
10.2.8 Dashboard 10-18
[Link] Visualizations Appearing Earlier are No Longer Available in Dashboard 10-18
[Link] Dashboard Layout Reset after You Resized/moved the Visualizations 10-18
[Link] Streaming Visualizations Do not Show Any Data 10-18
10.2.9 Live Output 10-18
[Link] Issues with Live Output 10-19
[Link] Missing Events due to Faulty Data 10-20
10.2.10 Pipeline Deployment Failure 10-21
xiv
1
Overview
Introduction
Key Features of GGSA
GGSA Architecture
Steps to build Continuous-ETL and Realtime-Analytics Pipelines
1.1 Introduction
The Oracle GoldenGate Stream Analytics (GGSA) runtime component is a complete solution
platform for building applications to filter, correlate, and process events in real-time. With
flexible deployment options of stand-alone Spark or Hadoop-YARN, it proves to be a versatile,
high-performance event processing engine. GGSA enables Fast Data and Internet of Things
(IOT) – delivering actionable insight and maximizing value on large volumes of high velocity
data from varied data sources in real-time. It enables distributed intelligence and low latency
responsiveness by pushing business logic to the network edge.
1-1
Chapter 1
GGSA Architecture
Acquiring data
Stream Analytics can acquire data from any of the following on-premises and cloud-native data
sources:
• GoldenGate: Natively integrated with Oracle GoldenGate, Stream Analytics offers data
replication for high-availability data environments, real-time data integration, and
transactional change data capture.
• Oracle Cloud Streaming: Ingest continuous, high-volume data streams that you can
consume or process in real-time.
• Kafka: A distributed streaming platform used for metrics collection and monitoring, log
aggregation, and so on.
• Java Message Service: Allows java-based applications to send, receive, and read
distributed communications.
Processing data
With Stream Analytics, you can filter, correlate, and process events in real-time.
Perform actions on the data
After Stream Analytics processes the data, you can output the results to any one of the
following external target data sources:
• Coherence
• Kafka
• Oracle Cloud Streaming
• Java Message Service
• Database
• Notification
• REST
1-2
Chapter 1
Steps to build Continuous-ETL and Realtime-Analytics Pipelines
1-3
2
Install
2-1
Chapter 2
Planning Your Installation
2-2
Chapter 2
Installing GoldenGate Stream Analytics
• Note:
Install Spark and JDK in the same node on which you plan to install Oracle
Stream Analytics. See Installing GoldenGate Stream Analytics.
2-3
Chapter 2
Configuring the Metadata Store
1. Create a directory, for example, spark-downloads, and download Apache Spark into the
newly created folder and as specified by versions, detailed in the Planning Your Installation
section.
2. Extract the Spark archive to a local directory.
You can see a subfolder, spark-*.*.*-bin-hadoop*.*.
3. Create a new directory, for example, OSA-19 and download OSA-[Link].*.zip from
Oracle eDelivery and extract it into the newly created folder.
You can find the OSA-[Link].*-[Link] file in the OSA-[Link].* zip file.
4. Extract the downloaded file. You should now see a subfolder OSA-[Link].*.
5. Review the file OSA-[Link].*-[Link] in the OSA-19 folder.
6. Set the environment variables:
<New id="osads"
class="[Link]">
<Arg>
<Ref refid="wac"/>
</Arg>
<Arg>jdbc/OSADataSource</Arg>
<Arg>
<New class="[Link]">
<Set
name="URL">jdbc:oracle:thin:@[Link]:OSADB</Set>
<Set name="User">OSA_USER</Set>
<Set name="Password">
<Call class="[Link]"
name="deobfuscate">
<Arg> OBF:OBFUSCATED_PASSWORD</Arg>
</Call>
</Set>
<Set name="connectionCachingEnabled">true</
Set>
2-4
Chapter 2
Configuring the Metadata Store
<Set name="connectionCacheProperties">
<New class="[Link]">
<Call name="setProperty"><Arg>MinLimit</Arg><Arg>1</Arg></
Call>
<Call name="setProperty"><Arg>MaxLimit</Arg><Arg>15</
Arg></Call>
<Call name="setProperty"><Arg>InitialLimit</Arg><Arg>1</
Arg></Call>
</New>
</Set>
</New>
</Arg>
</New>
3. Decide on an OSA schema username and a plain-text password. For illustration, say osa
as schema user name and alphago as password.
Change directory to top-level folder OSA-[Link].* and execute the following
command:java -cp ./lib/ [Link]
[Link] osa <your password>
For example, java -cp ./lib/ [Link]
[Link] osa alphago
You should see results like below on console:
2019-06-18 14:14:45.114:INFO::main: Logging initialized @1168ms to
[Link]
OBF:<obfuscated password>
MD5:34d0a556209df571d311b3f41c8200f3
CRYPT:osX/8jafUvLwA
4. Note down the obfuscated password string that is displayed (shown in bold), by copying it
to clipboard or notepad.
5. Change database host, port, SID, osa schema user name and osa schema password
fields marked in bold in the code in Step 2a.
Example - jdbc:oracle:thin:@[Link]:ORCL
SAMPLE [Link]
<?xml version="1.0"?>
<!DOCTYPE Configure PUBLIC "-//Jetty//Configure//EN" "[Link]
jetty/configure_9_3.dtd">
2-5
Chapter 2
Configuring the Metadata Store
<Set
name="URL">jdbc:oracle:thin:@[Link]:OSADB</Set>
<Set name="User">osa_prod</Set>
<Set name="Password">
<Call class="[Link]"
name="deobfuscate">
<Arg>OBF:1ggz1j1u1k8q1leq1v2h1w8v1v1x1lcs1k5g1iz01gez</Arg>
</Call>
</Set>
<Set name="connectionCachingEnabled">true</Set>
<Set name="connectionCacheProperties">
<New class="[Link]">
<Call name="setProperty"><Arg>MinLimit</Arg><Arg>1</
Arg></Call>
<Call name="setProperty"><Arg>MaxLimit</Arg><Arg>15</
Arg></Call>
<Call name="setProperty"><Arg>InitialLimit</
Arg><Arg>1</Arg></Call>
</New>
</Set>
</New>
</Arg>
</New>
2-6
Chapter 2
Configuring the Metadata Store
</Arg>
</New>
-->
<!-- SAMPLE OSA DATASOURCE CONFIGURATION FOR MYSQL-->
<!--
</Configure>
Note:
Do not use a hyphen in the OSA metadata username, in the [Link]
However, before running the above script, you must configure the datasource in the
datasource configuration file at${OSA_HOME}/osa-base/etc/jetty-osa-
[Link].
To configure ATP/ADW as metadata store, first comment the Oracle and MYSQL sections,
while uncommenting the ADW/APT section in [Link] file.
Below is the template for the datasource configuration for ATP/ADW database:
[Link]
<?xml version="1.0"?>
<!DOCTYPE Configure PUBLIC "-//Jetty//Configure//EN" "[Link]
2-7
Chapter 2
Initializing Metadata Store
jetty/configure_9_3.dtd">
<Configure id="Server" class="[Link]">
<New id="osads" class="[Link]">
<Arg>
<Ref refid="wac"/>
</Arg>
<Arg>jdbc/OSADataSource</Arg>
<Arg>
<New class="[Link]" type="adw">
<Set name="URL">jdbc:oracle:thin:@{service_name}?
TNS_ADMIN={wallet_absolute_path}</Set>
<Set name="User">{osa_db_user}</Set>
<Set name="Password">
<Call class="[Link]"
name="deobfuscate">
<Arg>{obfuscated_password}</Arg>
</Call>
</Set>
<Set name="connectionCachingEnabled">true</Set>
<Set name="connectionCacheProperties">
<New class="[Link]">
<Call name="setProperty"><Arg>MinLimit</Arg><Arg>1</
Arg></Call>
<Call name="setProperty"><Arg>MaxLimit</Arg><Arg>15</
Arg></Call>
<Call name="setProperty"><Arg>InitialLimit</
Arg><Arg>1</Arg></Call>
</New>
</Set>
</New>
</Arg>
</New>
</Configure>
Note:
In the above template, replace the variables in {} as below:
• {service_name} - one of the service names listed in the [Link] file inside
the wallet
• {wallet_absolute_path} - the absolute path of wallet folder on the machine where
OSA is installed
• {osa_db_user} - the username to create the osa metadata. This username and
schema will be created by the 'dbroot' user provided in above script.
• {obfuscated_password} - the Obfuscated password for {osa_db_user}
2-8
Chapter 2
Initializing Metadata Store
After installing GGSA, you need to configure the metadata store with the database admin
credential details and the version of GGSA as required.
To initialize the metadata store, you need database admin credentials with sysdba privileges:
1. Change directory to OSA-[Link].*/osa-base/bin.
2. Execute the following command:./[Link] dbroot=<db sys user>
dbroot_password=<db sys password> For example, ./[Link] dbroot=AlphaUser
dbroot_password=AlphaPassword
3. Following console messages will be displayed indicating the OSA schema was created
and prompting for password for osaadmin user.
The following console messages indicates that the GGSA schema is created and the
metadata store is successfully initialized:
4. Enter password:
5. Re-enter password:
6. Run ./[Link] to complete schema creation and metadata initialization.
7. If you don’t see the above messages, check the OSA-[Link].*/osa-base/logs folder to
identify the cause and potential solution.
Note:
If you do not have the database admin credentials, ask your database
administrator to create a GoldenGate Stream Analytics database user by using
the SQL scripts available in the OSA-[Link].*/osa-base/sql folder. The
GoldenGate Stream Analytics database username must match the one
configured in [Link].
2-9
Chapter 2
Jetty Properties File
Note:
It is recommended that you configure these properties at the installation stage, to
avoid restarting your server, if configured at a later stage.
Note:
If you do not specify explicitly the host header in your request, the default value is
host-server:port, where the OSA jetty server is running. Hence you must
specify the port number along with the server address.
• [Link]
You can restrict the x-forwarded-host header values to the values defined with this
property.
Example: [Link]= [Link], [Link],
localhost
Here the value of the x-forwarded-host header can be only of these three domains listed.
Commenting out this property with a # will allow all values for the header. If no domain is
entered, that is, if the value of the property is empty, then this header is not supported.
• [Link]
A comma separated list of response headers, which will be sent along with response for
every request.
Example: [Link]="x-frame-options: sameorigin, X-Content-Type-
Options: nosniff"
By default the above 2 response headers are set.
– x-frame-options: sameorigin will prevent clickjack attacking.
2-10
Chapter 2
Adjusting Jetty Threadpool
<New id="threadPool"
class="[Link]">
<Set name="minThreads" type="int"><Property
name="[Link]"
deprecated="[Link]" default="100"/></Set>
<Set name="maxThreads" type="int"><Property name="[Link]"
deprecated="[Link]" default="2000"/></Set>
<Set name="reservedThreads" type="int"><Property
name="[Link]" default="-1"/></Set>
<Set name="idleTimeout" type="int"><Property
name="[Link]" deprecated="[Link]"
default="60000"/></Set>
<Set name="detailedDump" type="boolean"><Property
name="[Link]" default="false"/></Set>
</New>
</Configure>
Note:
Install Oracle GoldenGate Big Data on the same machine and with the same
user as OSA.
2-11
Chapter 2
Maven Setting for GoldenGate Big Data Handlers
to
OSA_HOME="$( cd "$(dirname "../../../")" >/dev/null 2>&1 ; pwd -P )"
Note:
Update the maven home path before initialization of the metadata store, or you will
have to restart GGSA after this update.
<proxy>
<id>optional</id>
<active>true</active>
<protocol>http</protocol>
<username>proxyuser</username>
<password>proxypass</password>
<host>[Link]</host>
<port>80</port>
<nonProxyHosts>[Link]|[Link]</nonProxyHosts>
</proxy>
Note:
Username and password field is required if the proxy is protected.
2-12
Chapter 2
GoldenGate Stream Analytics Hardware Requirements for Enterprise Deployment
Note:
Update the [Link] before initialization of the metadata store, or you will
have to restart GGSA after this update.
Note:
The two-node Kafka cluster can be avoided if customer already has a Kafka cluster in
place and is fine with OSA leveraging that cluster for its internal usage.
Based on the above estimates, total cores for design-tier is 12 and approximate memory is 112
GB RAM. Jetty instances can be independently scaled as the number of users increase.
2-13
Chapter 2
GoldenGate Stream Analytics Hardware Requirements for Enterprise Deployment
Data Tier
The deployed pipelines are run on the YARN or Spark cluster. You can use existing YARN/
Spark clusters if you have sufficient spare capacity.
Sizing Guidelines
Use the following sizing guidelines to run GGSA pipelines. Ensure that the pipelines are
deployed on shared storage, so that the pipeline code and libraries are accessible from all
nodes in the YARN/Spark cluster. GGSA supports NFS for shared storage but if you want to
use HDFS, the hardware needs two more nodes.
• – 2 nodes with 4+ cores, 16+GB RAM, and 500 GB local disk to run HDFS cluster, two
instances of HDFS name and data nodes.
The Spark tier is where work happens and the Spark cluster size depends on
• Number of pipelines that will simultaneously run
• Logic in each pipeline
• Desired degree of parallelism
For each streaming pipeline the number of cores and memory gets computed based on a
required degree of parallelism. As an example, consider a pipeline ingesting data from
customer’s Kafka topic T with 3 partitions using direct ingestion. Direct ingestion is where no
Spark Receivers are used. In this case, the minimum number of processes that you need to
run for optimal performance is as follows: 1 Spark Driver Process + 3 Executor processes, 1
for each Kafka Topic partition. Each Executor process needs a minimum of 2 cores.
The number of cores for a pipeline can be computed as
--executor-cores = 1 + Number of Executors * 2
2-14
Chapter 2
GoldenGate Stream Analytics Hardware Requirements for Enterprise Deployment
If you are considering GGSA for POCs and not production, then you can use the
following configuration:
Design Tier
• An instance of the Jetty running on a 4+ core node with a 32+ GB of RAM.
• An instance of MySQL/Oracle for metadata store on a 4+ core node with a 16+ GB of
RAM.
• A node of the Kafka cluster running on a 4+ core node with 16+ GB of RAM.
Note:
This is a separate Kafka cluster for GGSA’s internal use and for interactively
designing pipelines.
Data Tier
• A Hadoop Distributed File System (HDFS) cluster node running on 4+ core physical node
with 16+ GB of RAM.
• 2 nodes of the YARN/Spark cluster each running on a 4+ core physical node with a 16+
GB of RAM.
Development Mode Configurations
Design Tier
2-15
Chapter 2
Retaining https and Disabling http
• 1 node with 4+ cores and 16+ GB of RAM for 1 instance of Jetty, 1 instance of MySQL DB,
and 1 instance of Kafka+ZooKeeper
Data Tier
• 1 node with 4+ cores and 16+ GB of RAM for 1 instance of HDFS and 1 instance of YARN/
Spark.
Note:
The password is a plain-text password.
3. Click the user name at the top right corner of the screen.
4. Click System Settings.
5. Click Environment.
6. Select the Runtime Server. See the sections below for Yarn and Spark Standalone
runtime configuration details.
Yarn Configuration
1. • YARN Resource Manager URL: Enter the URL where the YARN Resource Manager
is configured.
2-16
Chapter 2
Setting up Runtime for GoldenGate Stream Analytics Server
• Storage: Select the storage type for pipelines. To submit a GGSA pipeline to Spark,
the pipeline has to be copied to a storage location that is accessible by all Spark
nodes.
– If the storage type is WebHDFS:
* Path: Enter the WebHDFS directory (hostname:port/path), where the
generated Spark pipeline will be copied to and then submitted from. This
location must be accessible by all Spark nodes. The user specified in the
authentication section below must have read-write access to this directory.
* HA Namenodes: Set the HA namenodes. If the hostname in the above URL
refers to a logical HA cluster, specify the actual namenodes here, in the
format:Hostname1:Port, Hostname2:Port.
– If storage type is HDFS:
* Path: The path could be <HostOrIPOfNameNode><HDFS Path>. For
example, [Link]/user/oracle/ggsapipelines. Hadoop
user must have Write permissions. The folder will automatically be created if it
does not exist.
* HA Namenodes: If the hostname in the above URL refers to a logical HA
cluster, specify the actual namenodes here, in the format:Hostname1:Port,
Hostname2:Port.
– If storage type is NFS:
Path: The path could be /oracle/spark-deploy.
Note:
/oracle should exist and spark-deploy will automatically be created if it
does not exist. You will need Write permissions on the /oracle directory.
2. Hadoop Authentication:
• Simple authentication credentials:
– Protection Policy: Select a protection policy from the drop-down list. This value
should match the value on the cluster.
– Username: Enter the user account to use for submitting Spark pipelines. This user
must have read-write access to the Path specified above.
• Kerberos authentication credentials:
– Protection Policy: Select a protection policy from the drop-down list. This value
should match the value on the cluster.
– Kerberos Realm: Enter the domain on which Kerberos authenticates a user, host,
or service. This value is in the [Link] file.
– Kerberos KDC: Enter the server on which the Key Distribution Center is running.
This value is in the [Link] file.
– Principal: Enter the GGSA service principal that is used to authenticate the GGSA
web application against Hadoop cluster, for application deployment. This user
should be the owner of the folder used to deploy the GGSA application in HDFS.
You have to create this user in the yarn node manager as well.
– Keytab: Enter the keytab pertaining to GGSA service principal.
2-17
Chapter 2
Setting up Runtime for GoldenGate Stream Analytics Server
– Yarn Resource Manager Principal: Enter the yarn principal. When Hadoop
cluster is configured with Kerberos, principals for hadoop services like hdfs, https,
and yarn are created as well.
3. Yarn master console port: Enter the port on which the Yarn master console runs. The
default port is 8088.
4. Click Save.
Spark Standalone
1. Select the Runtime Server as Spark Standalone, and enter the following details:
• Spark REST URL: Enter the Spark standalone REST URL. If Spark standalone is HA
enabled, then you can enter comma-separated list of active and stand-by nodes.
• Storage: Select the storage type for pipelines. To submit a GGSA pipeline to Spark,
the pipeline has to be copied to a storage location that is accessible by all Spark
nodes.
– If the storage type is WebHDFS:
* Path: Enter the WebHDFS directory (hostname:port/path), where the
generated Spark pipeline will be copied to and then submitted from. This
location must be accessible by all Spark nodes. The user specified in the
authentication section below must have read-write access to this directory.
* HA Namenodes: If the hostname in the above URL refers to a logical HA
cluster, specify the actual namenodes here, in the format:Hostname1:Port,
Hostname2:Port.
– If storage type is HDFS:
* Path: The path could be <HostOrIPOfNameNode><HDFS Path>. For
example, [Link]/user/oracle/ggsapipelines. Hadoop
user must have Write permissions. The folder will automatically be created if it
does not exist.
* HA Namenodes: If the hostname in the above URL refers to a logical HA
cluster, specify the actual namenodes here, in the format:Hostname1:Port,
Hostname2:Port.
This field is applicable only when the storage type is HDFS.
– Hadoop Authentication for WebHDFS and HDFS Storage Types:
* Simple authentication credentials:
* Protection Policy: Select a protection policy from the drop-down list.
* Username: Enter the user account to use for submitting Spark pipelines.
This user must have read-write access to the Path specified above.
* Kerberos authentication credentials:
* Protection Policy: Select a protection policy from the drop-down list.
* Kerberos Realm: Enter the domain on which Kerberos authenticates a
user, host, or service. This value is in the [Link] file.
* Kerberos KDC: Enter the server on which the Key Distribution Center is
running. This value is in the [Link] file.
* Principal: Enter the GGSA service principal that is used to authenticate
the GGSA web application against Hadoop cluster, for application
deployment. This user should be the owner of the folder used to deploy
2-18
Chapter 2
Validating Data Flow to GoldenGate Stream Analytics
the GGSA application in HDFS. You have to create this user in the yarn
node manager as well.
* Keytab: Enter the keytab pertaining to GGSA service principal.
* Yarn Resource Manager Principal: Enter the yarn principal. When
Hadoop cluster is configured with Kerberos, principals for hadoop services
like hdfs, https, and yarn are created as well.
– If storage type is NFS:
Path: The path could be /oracle/spark-deploy.
Note:
/oracle should exist and spark-deploy will automatically be created if it
does not exist. You will need Write permissions on the /oracle directory.
2. Spark standalone master console port: Enter the port on which the Spark standalone
console runs. The default port is 8080.
Note:
The order of the comma-separated ports should match the order of the comma-
separated spark REST URLs mentioned in the Path.
Note:
You can change your Spark standalone server username and password in this
screen. The username and password fields are left blank, by default.
5. Click Save.
ProductLn,ProductType,Product,OrderMethod,CountrySold,QuantitySold,UnitSale
Price
Personal Accessories,Watches,Legend,Special,Brazil,1,240
Outdoor Protection,First Aid,Aloe Relief,E-mail,United States,3,5.23
Camping Equipment,Lanterns,Flicker Lantern,Telephone,Italy,3,35.09
Camping Equipment,Lanterns,Flicker Lantern,Fax,United States,4,35.09
Golf Equipment,Irons,Hailstorm Steel Irons,Telephone,Spain,5,461
2-19
Chapter 2
Terminating GoldenGate Stream Analytics
2. In the Catalog, as shown in the image below, click Create New Item, and then click
Stream. create a stream of type File.
3. In the Type Properties page of the Create Stream dialog box, provide the Name,
Description, and Tags for the Stream, select the Stream Type as File, and then select
Create Pipeline with this source (Launch Pipeline Editor).
4. Click the Next button to navigate to the Source Details page of the Create Stream dialog
box.
5. In the Source Details page, click Upload file to upload the [Link] file, and then click
Next to navigate to the Data Format page.
6. In the Data Format page, select the CSV Predefined Format as Default and select the
First record as header, and then click Next to navigate to the Shape page.
7. In the Shape page, verify that the shape of the event is successfully inferred as in the
following image, and then click Save.
8. In the Create Pipeline dialog box, enter the Name, Description, Tags of the pipeline,
select the Stream that you created, and then click Save:
You can see the pipeline editor and you can see the message Starting Pipeline followed
by the message Listening to Events.
Note:
This is the first access of the cluster and it takes time to copy libraries, please be
patient. You should eventually see the screenshot below with single node
representing the stream source.
2-20
Chapter 2
Upgrading GoldenGate Stream Analytics
Note:
You can skip this step, if you are upgrading to GGSA version [Link].8.
2-21
3
Configure
Configure Runtime Environment
Configure Users
3-1
Chapter 3
Configure Runtime Environment
Note:
Enter the tenancyName and userName, not tenancy OCID and user
OCID. Similarly, enter the stream pool ID and not the stream pool name.
You can retrieve this information from the OCI console. This field is enabled only if
you have checked the SASL option.
– Password: Enter the SASL password, which is an authentication token that you
can generate on the User Details page, of the OCI console.
Note:
Copy the authentication token when you create it, and save it for future
use. You can not retrieve it at a later stage.
Kafka Topics
3-2
Chapter 3
Configure Runtime Environment
Group IDs
3-3
Chapter 3
Configure Runtime Environment
Note:
/oracle should exist and spark-deploy will automatically be created if it
does not exist. You will need Write permissions on the /oracle directory.
5. Spark standalone master console port: Enter the port on which the Spark standalone
console runs. The default port is 8080.
Note:
The order of the comma-separated ports should match the order of the comma-
separated spark REST URLs mentioned in the Path.
3-4
Chapter 3
Configure Runtime Environment
Note:
You can change your Spark standalone server username and password in this
screen. The username and password fields are left blank, by default.
8. Click Save.
Note:
/oracle should exist and spark-deploy will automatically be created if it
does not exist. You will need Write permissions on the /oracle directory.
5. Hadoop Authentication:
• Simple authentication credentials:
– Protection Policy: Select a protection policy from the drop-down list. This value
should match the value on the cluster.
3-5
Chapter 3
Configure Runtime Environment
– Username: Enter the user account to use for submitting Spark pipelines. This user
must have read-write access to the Path specified above.
• Kerberos authentication credentials:
– Protection Policy: Select a protection policy from the drop-down list. This value
should match the value on the cluster.
– Kerberos Realm: Enter the domain on which Kerberos authenticates a user, host,
or service. This value is in the [Link] file.
– Kerberos KDC: Enter the server on which the Key Distribution Center is running.
This value is in the [Link] file.
– Principal: Enter the GGSA service principal that is used to authenticate the GGSA
web application against Hadoop cluster, for application deployment. This user
should be the owner of the folder used to deploy the GGSA application in HDFS.
You have to create this user in the yarn node manager as well.
– Keytab: Enter the keytab pertaining to GGSA service principal.
– Yarn Resource Manager Principal: Enter the yarn principal. When Hadoop
cluster is configured with Kerberos, principals for hadoop services like hdfs, https,
and yarn are created as well.
6. Yarn master console port: Enter the port on which the Yarn master console runs. The
default port is 8088.
7. Click Save.
3-6
Chapter 3
Configure Runtime Environment
The default security list for the Regional Subnet must allow bidirectional traffic to edge/OSA
node so create a stateful rule for destination port 443. Also create a similar Ingress rule for port
7183 to access the BDS Cloudera Manager via the OSA edge node. An example Ingress rule
is shown below.
You will also be able to access your Cloudera Manager using the same public IP of GGSA
instance by following steps below. Please note this is [Link] to GGSA box and run the
following port forward commands so you can access the Cloudera Manager via the GGSA
instance:
sudo firewall-cmd --add-forward-port=port=7183:proto=tcp:toaddr=<IP
address of Utility Node running the Cloudera Manager console>
sudo firewall-cmd --runtime-to-permanent
sudo sysctl net.ipv4.ip_forward=1
3-7
Chapter 3
Configure Runtime Environment
You should now be able to access the Cloudera Manager using the URL [Link]
IP of GGSA>:7183.
[Link].3.2 Prerequisites
1. Retrieve IP addresses of BDS cluster nodes from OCI console as shown in screenshot
below. Alternatively, you can get the FQDN for BDS nodes from the Cloudera Manager as
shown below.
2. Reconfigure YARN virtual cores using Cloudera Manager as shown below. This will allow
many pipelines to run in the cluster and not be bound by actual physical cores.
3-8
Chapter 3
Configure Runtime Environment
• Container Virtual CPU Cores: This is the total virtual CPU cores available to YARN
Node Manager for allocation to Containers. Please note this is not limited by physical
cores and you can set this to a high number, say 32 even for VM standard 2.1.
• Container Virtual CPU Cores Minimum: This is the minimum vcores that will be
allocated by YARN scheduler to a Container. Please set this to 2 since CQL engine is
a long-running task and will require a dedicated vcore.
Container Virtual CPU Cores Maximum: This is the maximum vcores that will be
allocated by YARN scheduler to a Container. Please set this to a number higher than 2
say 4.
Note:
This change will require a restart of the YARN cluster from Cloudera Manager.
3-9
Chapter 3
Configure Runtime Environment
6. Set HA Namenode to Private IPs or hostnames of all master nodes (comma separated),
starting with the one running active NameNode server. In the next version of GGSA, the
ordering will not be needed. For example,
[Link],
[Link].
7. Set Yarn Master Console port to 8088 or as configured in BDS.
8. Set Hadoop authentication to Kerberos.
9. Set protection policy to privacy. Please note this should match the value in HDFS
configuration property [Link].
10. Set Kerberos Realm to [Link].
11. Set Kerberos KDC to private IP or hostname of BDS master node 0. For example,
[Link].
12. Set principal to bds@[Link]. See this documentation to
create a Kerberos principal (e.g. bds) and add it to hadoop admin group, starting with step
Connect to Cluster's First Master Node and through the step Update HDFS
Supergroup.
Note:
You can hop/ssh to the master node using your GGSA node as the Bastion. You
will need your ssh private key to be available on GGSA node though. Restart
your BDS cluster as instructed in the documentation.
[opc@bdsggsa ~]$ ssh -i id_rsa_private_key
opc@[Link]
13. Make sure the newly created principal is added to Kerberos keytab file on the master node
as shown:
bdsmn0 # sudo [Link]
[Link]: ktadd -k /etc/[Link] bds@[Link]
14. Fetch the keytab file using sftp and set Keytab field in system settings by uploading the
same.
3-10
Chapter 3
Configure Runtime Environment
15. Set Yarn Resource Manager principal which should be in the format yarn/<FQDN of
BDS MasterNode running Active Resource Manager>@KerberosRealm. For
example, yarn/
[Link]@BDACLOUDSERVI
[Link].
Sample System Settings Screen:
3-11
Chapter 3
Configure Runtime Environment
6. Set HA Namenode to Private IP or Hostname of the BDS Master node. For example,
[Link].
7. Set Yarn Master Console port to 8088 or as configured in BDS
8. Set Hadoop authentication to Simple and leave Hadoop protection policy at authentication
if available
9. Set username to oracle.
10. Click Save.
3-12
Chapter 3
Configure Runtime Environment
• Log Level: Select a log level for unpublished pipelines, from the drop-down list.
Note:
Reset the default log level of draft pipelines to WARNING, and of published
pipelines, to ERROR.
Note:
When you publish the pipeline for the first time, the input stream is read
based on the offset value you have selected in this drop-down list. On a
subsequent publish, the value you have selected here is not considered, and
the input stream is read from where it was last left off.
• Reset Offset: Select this option to read the input stream based on the offset value
selected in the Input Topics Offset drop-down list.
Note:
If you are using two Kafka streams as an input to the pipeline, the offset is
not preserved and the pipeline starts from the current timestamp. With a
single stream the offset is maintained and the pipeline can read from the
previous state of it.
3. Click Save.
3-13
Chapter 3
Configure Runtime Environment
• No proxy for: Set a list of hosts that should be reached directly, bypassing the proxy.
This is a list of patterns separated by the delimiter |. The patterns can start or end with
a * for wildcards.
3. Click Save.
3-14
Chapter 3
Configure Runtime Environment
In this example all the tables with names starting with DEMO will be displayed when
using references or targets.
• Oracle DB Geofence: Enter the query for geofence tables and columns details.
Enter a SQL query to fetch the table column details, for all the tables to be used for
creating Geofences. The SELECT query must contain the following attributes for each
column: table_name, column_name, data_type, data_length, data_precision,
data_scale. The tables listed in for geofence must have at least one column of the
SDO_GEOMETRY type. Hence the following where clause is mandatory: WHERE
table_name IN (SELECT TABLE_NAME FROM all_tab_columns WHERE
data_type='SDO_GEOMETRY'.
Example: SELECT table_name, column_name, data_type, data_length,
data_precision, data_scale FROM all_tab_columns WHERE table_name
IN (SELECT TABLE_NAME FROM all_tab_columns WHERE
data_type='SDO_GEOMETRY') AND table_name LIKE GEO ORDER BY
table_name desc. In this example, all the tables with names starting with GEO will
be listed in the descending order of the table name.
3. Click Save.
3-15
Chapter 3
Configure Users
Note:
You can change [Link] to daily, hourly,
minutely. This is to enable log rollover based on time.
+----+----------+--------------------------------------+
| id | username | pwd |
+----+----------+--------------------------------------+
| 1 | osaadmin | MD5:201f00b5ca5d65a1c118e5e32431514c |
+----+----------+--------------------------------------+
3-16
Chapter 3
Configure Users
where osaadmin is the pre-configured user along with the encrypted password.
When you execute a query to pull in all the data from the osa_user_roles table, you can see
the following:
+---------+---------+
| user_id | role_id |
+---------+---------+
| 1 | 1 |
+---------+---------+
Note:
Replace *.* with the current version of GGSA.
3-17
Chapter 3
Configure Users
where NewUser is the name of the user and <password> is the password that you want to
obfuscate or encrypt.
You will see a similar screen on your terminal:
For more information about running the password utility, see Configuring Secure
Password.
3. Connect to the database using the database user credentials that you have configured in /
osa-base/etc/[Link].
4. Insert a record into the osa_users table using any one of the following commands:
or
or
5. Insert a record into the osa_user_roles table using the following command:
Important:
Currently, Oracle GoldenGate Stream Analytics supports only one user role, i.e,
the administrator role. So the role_id value must always be 1.
You can now login to Oracle GoldenGate Stream Analytics as NewUser using <password>.
Repeat these steps to create as many users as you require.
3-18
Chapter 3
Configure Users
3-19
Chapter 3
Configure Users
1. Execute the following command from SQLPLUS or SQLDeveloper tools to remove a user:
This command deletes the user with the id value as 2, i.e, the second user in the database.
2. Execute the following command to delete the user role corresponding to the user in the
above step:
osa-demo-LDAP {
[Link] required
debug="true"
contextFactory="[Link]"
hostname=<hostname> <!-- hostname of LDAP -->
port="389"
authenticationMethod="simple"
forceBindingLogin="true"
userBaseDn="l=emea,dc=oracle,dc=com"
userRdnAttribute="uid"
3-20
Chapter 3
Configure Users
userIdAttribute="mail"
userPasswordAttribute="userPassword"
userObjectClass="person"
roleBaseDn="l=emea,dc=oracle,dc=com"
roleNameAttribute="opn_access_level"
roleMemberAttribute="targetdn"
roleObjectClass="person";
};
osa-demo-LDAP {
[Link] required
debug="true"
contextFactory="[Link]"
hostname=<hostname> <!-- hostname of LDAP -->
port="389"
authenticationMethod="simple"
forceBindingLogin="true"
userBaseDn="l=amer,dc=oracle,dc=com"
userRdnAttribute="uid"
userIdAttribute="mail"
userPasswordAttribute="userPassword"
userObjectClass="person"
roleBaseDn="l=amer,dc=oracle,dc=com"
roleNameAttribute="employeetype"
roleMemberAttribute="targetdn"
roleObjectClass="organizationalPerson";
};
userBaseDn="l=amer,dc=oracle,dc=com"
roleBaseDn="l=amer,dc=oracle,dc=com"
If in Asia Pacific:
userBaseDn="l=apac,dc=oracle,dc=com"
roleBaseDn="l=apac,dc=oracle,dc=com"
If in Europe:
userBaseDn="l=emea,dc=oracle,dc=com"
roleBaseDn="l=emea,dc=oracle,dc=com"
3-21
Chapter 3
Configure Users
1. Ensure that role name is updated in the [Link] file located at /osa-base/etc/
[Link]:
<auth-constraint>
<role-name>developer</role-name>
</auth-constraint>
osa_demo_ldap {
[Link] required
debug="true"
contextFactory="[Link]"
hostname=<hostname> <!-- this is the active directory server
hostname -->
port="389" <!-- this is the active directory server port -->
bindDn="CN=Administrator,CN=Users,DC=corp,DC=oradev,DC=com"
bindPassword=<password> <!-- If the active directory server
allows anonymous login, no need to provide bindDn and bindPassword. Else,
set the active directory server admin DN and password -->
authenticationMethod="simple" <!-- if the active directory
server allows anonymous login then set to 'none' otherwise set it to
'simple'-->
forceBindingLogin="true"
userBaseDn="l=amer,dc=oracle,dc=com" <!-- user attributes as
per user setup in active directory server -->
userRdnAttribute="uid" <!-- user attributes as per user setup
in active directory server -->
userIdAttribute="mail" <!-- user attributes as per user setup
in active directory server -->
userPasswordAttribute="userPassword" <!-- user attributes as
per user setup in active directory server -->
userObjectClass="person" <!-- user attributes as per user
setup in active directory server -->
roleBaseDn="l=amer,dc=oracle,dc=com" <!-- role (group)
attributes as per user setup in active directory server -->
roleNameAttribute="opn_access_level" <!-- role (group)
attributes as per user setup in active directory server -->
roleMemberAttribute="targetdn" <!-- role (group) attributes as
per user setup in active directory server -->
roleObjectClass="person"; <!-- role (group) attributes as per
user setup in active directory server -->
};
3-22
Chapter 3
Configure Users
3-23
4
Manage
Manage Connections
Manage Streams
Manage References
Manage Targets
Manage GG Change Data Stream
Embedded Ignite Cache
Manage Pipelines
4.1 Connections
Create Connections
Manage Connections
4-1
Chapter 4
Connections
4-2
Chapter 4
Connections
Note:
Special characters in the password are treated as wildcard patterns, causing
the connection to fail . If your password contains special characters, enclose
the password within double quotes (the double quotes with ASCII Value 34).
6. Click Test Connection, to ensure that you have successfully created a connection.
7. Click Save.
4-3
Chapter 4
Connections
4-4
Chapter 4
Connections
4-5
Chapter 4
Connections
• Service Manager Host: Enter the name or IP address of the GoldenGate Service
Manager.
• Service Manager Port: Enter the port on the Service Manager, to connect to
GoldenGate.
Set the port to 443 to create a connection to the Goldengate instance running on an
OCI Goldengate stack, because you can access only port 443, by default.
• GG Username: Enter the username to authenticate the GoldenGate connection.
• GG Password: Enter the password for the GoldenGate connection.
• Is SSL?: Select this option if the Goldengate instance uses a SSL based connection
• Is GG Marketplace?: Select this option if the GG instance is running on OCI
Marketplace Goldengate stack.
6. Click Test Connection, to ensure that you have successfully created a connection.
7. Click Save.
4-6
Chapter 4
Connections
Note:
Retain the [Link] and [Link] file names exactly as they are.
7. Click Save.
4-7
Chapter 4
Connections
• Name: Enter a unique name for the connection. This is a mandatory field.
• Display Name: Enter a display name for the connection. If left blank, the Name field
value is copied.
• Description
• Tags
• Connection Type: The selected connection is displayed.
5. Click Next.
6. On the Connection Details screen, enter the following details:
• Core Site XML: Upload the [Link] file to connect to HDFS, where the files to
create external hive table are loaded. This is a mandatory field..
• Hdfs Site XML: Upload the [Link] file to connect to HDFS, where the files to
create external hive table are loaded. This is an optional field.
• Use Kerberos: In case of a kerberized cluster, you can provide the Kerberos principal
and keytab by enabling this option.
• Kerberos KDC: Provide the host having the Key Distribution Center(KDC).
• Kerberos Realm: Provide the kerberos realm.
• Kerberos Principal: Set the Kerberos principal for the hive service.
• Kerberos KeyTab: Upload the kerberos keytab file for the hive service.
• Hive JDBC URL: You can connect to hive database using following jdbc url :
jdbc:hive2://host:port/<DB_NAME>. The default port is 10000 and default database is
"default".
• Hive JDBC Username: Enter the username to connect to hive database. If you are
using the default database, this field can be left blank.
• Hive JDBC Password: Password used to connect to hive database. If you are using
the default database, this field can be left blank.
7. Click Test Connection, to ensure that you have successfully created a connection, and to
download the third-party libraries required to connect to hive database to create an
external table.
8. Click Save.
4-8
Chapter 4
Connections
4-9
Chapter 4
Connections
• Tags
• Connection Type: The selected connection is displayed.
4. Click Next.
5. On the Connection Details screen, enter the following details:
• Use Bootstrap: Check this box to use a bootstrap based connection.
• Zookeepers: Enter the zookeeper URL. Use this option only if you did not select the
Use Bootstrap box in the previous step.
• Kafka bootstrap: Enter the bootstrap URL.
• SSL: Check this box to connect to an SSL enabled Kafka cluster.
– Truststore Location: Locate and upload the truststore file. This field is applicable
only to connect to an SSL enabled Kafka cluster.
– Truststore Password: Enter the truststore password.
• SASL: Check this box if Kafka broker requires authentication.
– Security Protocol: From the drop-down, select the security protocol to be
associated with the Kafka connection.
– Security Mechanism: From the drop-down, select the security mechanism to be
associated with the Kafka connection.
– User Name: Enter the SASL username for the Kafka broker.
You can retrieve this information from the OCI console. This field is enabled only if
you have checked the SASL option.
– Password: Enter the SASL password, which is an authentication token that you
can generate on the User Details page, of the OCI console.
Note:
Copy the authentication token when you create it, and save it for future
use. You can not retrieve it at a later stage.
• MTLS: Select MLTS to enable 2-way authentication of both the user and the Kafka
broker.
– Truststore: Locate and upload the truststore file. This field is applicable only to
connect to an SSL enabled Kafka cluster.
– Truststore Password: Enter the truststore password.
– Keystore: Locate and upload the keystore file. This field is applicable only to
connect to an SSL enabled Kafka cluster.
– Keystore Password: Enter the keystore password.
6. Click Test Connection, to ensure that you have successfully created a connection.
7. Click Save.
4-10
Chapter 4
Connections
2. Hover the mouse over Connection and select HDFS from the submenu.
3. On the Type Properties screen, enter the following details:
• Name: Enter a unique name for the connection. This is a mandatory field.
• Display Name: Enter a display name for the connection. If left blank, the Name field
value is copied.
• Description
• Tags
• Connection Type: The selected connection is displayed.
4. Click Next.
5. On the Connection Details screen, enter the following details:
• [Link]: Upload the [Link] file with [Link], [Link],
[Link], and [Link]
properties.
Note:
Download the files [Link] and [Link] the Azure
website.
Note:
Retain the [Link] and [Link] file names exactly as they are.
7. Click Save.
4-11
Chapter 4
Connections
• Display Name: Enter a display name for the connection. If left blank, the Name field
value is copied.
• Description
• Tags
• Connection Type: The selected connection is displayed.
4. Click Next.
5. On the Connection Details screen, enter the following details:
• Connection Mode: Select the connection mode from the drop-down list:
– Server Address List:
* Server Address List: Connection to a list of Replicat set members or mongos.
This field accepts a comma separated list of hostnames:port. For example,
localhost1:27017, localhost2:27018.
* Authentication Mechanism: Select the authentication mechanism for the
connection, from the drop-down list. This is an optional field. Enter the
following details for the mechanism you select:
* Username: Enter the database account user name.
* Password: Enter the password for your database account.
* Credentials Source: Enter the source of the authentication
credentials, typically the database that the credentials have been
created in.
* Write Concern: Enter the value in JSON format. Accepted keys are w
and wtimeout. For example, {"w": "value" , "wtimeout":
"number"}
– Client URI:
* Client URI: Set the client URI in the format: mongodb://
[username:password@]host1[:port1][,host2[:port2],...
[,hostN[:portN]]][/[database][?options]]
* Authentication Mechanism:
* None: Select this option to disable connection authentication.
* SSL Server Certificate Validation:
* Trust Store File: Upload the truststore file.
* Trust Store Password: Set the truststore password.
* SSL Server/Client Certificate Validation:
* Trust Store File: Upload the truststore file.
* Trust Store Password: Set the truststore password.
* Key Store File: Upload the client certificate, for a two-way SSL
communication.
* Key Store Password: Set the keystore password.
6. Click Test Connection, to ensure that you have successfully created a connection.
7. Click Save.
4-12
Chapter 4
Connections
4-13
Chapter 4
Connections
• OCI Fingerprint : Enter the fingerprint of the API public key file that you uploaded to
OCI. For example, oci_api_key_public.pem.
• OCI Key File: Select the API private key file for signing API calls. For example,
oci_api_key.pem.
• Key Passphrase: Enter the Passphrase for the API private key file. The API key can
be passphrase-protected.
• OCI Tenancy OCID: Enter the OCID of the tenant in which the Object Store bucket is
defined.
• OCI Profile: Set the OCI profile. Default value is DEFAULT.
• OCI Namespace: Enter the OCI namespace that spans all compartments within a
region.
• Region: Enter the region in which tenancy is created. For a list of OCI regions, refer to
the Region Identifier column in the Regions and Availability Domains documentation.
• OCI Compartment OCID: Enter the OCID of the compartment in which the ONS topic
or Object Store is defined.
6. Click Test Connection, to validate the credentials to connect to OCI, and to ensure that
the dependent client libraries are downloaded from maven central repository.
7. Click Save.
4-14
Chapter 4
Connections
Note:
To integrate with OCI Notification service, you have to define an OCI Notification
service topic. Once you have defined a topic, note down the following parameters
that are required to send Messages from GGSA to OCI Notification:
• My Message: The message that target pushed to OCI Notification service topic.
• Topic: The topic created on OCI Notification service. The OSA target can publish
message to this topic.
• Email, Function, HTTPS, Slack: The subscriptions to the topic. All users who
have subscribed to the topic receive the message.
4-15
Chapter 4
Connections
4-16
Chapter 4
Connections
• Tags
• Connection Type: The selected connection is displayed.
4. Click Next.
5. On the Connection Details screen, enter the following details:
• Use Bootstrap: Check this box to use a bootstrap based connection. This is
mandatory for OSS connections. Connection to OSS can only be established using the
Bootstrap server option.
• Kafka bootstrap: Enter the bootstrap URL.
• SSL: Do not check this box when connecting to OCI Streaming Service.
• SASL: Check this box if Kafka broker requires authentication. This is mandatory for
OSS connections.
• User Name: Enter the SASL username for the Kafka broker, in the following format:
tenancyName/username/stream pool id
Note:
Enter the tenancyName and userName, not tenancy OCID and user OCID.
Similarly, enter the stream pool ID and not the stream pool name. Ensure
that the auto create topic is enabled for the stream pool ID.
You can retrieve this information from the OCI console. This field is enabled only if you
have checked the SASL option.
• Password — Enter the SASL password, which is an authentication token that you can
generate on the User Details page, of the OCI console.
Note:
Copy the authentication token when you create it, and save it for future use.
You can not retrieve it at a later stage.
6. Click Test Connection, to ensure that you have successfully created a connection.
7. Click Save.
Deleting a Connection
To delete a connection:
4-17
Chapter 4
Streams
1. Go to the Catalog page and hover the mouse over the connection that you want to delete.
2. Click the delete icon that appears to your right side on the screen.
3. On the Delete Confirmation screen, click Delete.
4.2 Streams
Create Streams
Manage Streams
Note:
Use File stream only for POCs and quick prototyping
• Read whole content: Select this option to read all the records in the file, at once. If
you uncheck this option, the engine reads one record at a time.
• Number of events per batch: Enter the number of records that you want to process
per batch. The default value is one, but you can specify the number of records to
process in each read. You can use this option only when Read Whole Content is
unchecked.
• Loop: Select this option to process the file in a loop.
• Data Format: Select CSV or JSON as the data format.
6. Click Next.
4-18
Chapter 4
Streams
7. On the Data Format screen, set the attributes for the selected the data format.
• For JSON data format:
– Allow Missing Column Names: Select this option to allow an input stream that
has a column undefined in the shape.
– Array in Multi-lines: Select this option to allow multi-line data formatting.
• For CSV data format:
– CSV Predefined Format: Select one of the predefined data format from the drop-
down list. For more information, see Predefined CSV Data Formats.
– First record as header: Select this option to use the first record as the header
row.
8. Click Next.
9. On the Shape screen, select one of the methods to define the shape:
• Infer Shape : Select this option to detect the shape automatically from the input data
stream.
• Select Existing Shape: Select one of the existing shapes from the drop-down list.
• Manual Shape : Select this option to infer the fields from a stream or file. You can also
update the datatype of the fields.
Note:
– To retrieve the entire JSON payload, add a new field with path $.
– To retrieve the content of the array, add a new field with path $
[arrayField].
In both the cases, the value returned is Text.
• From File: Select this option to infer the shape from a JSON schema file, or a JSON or
CSV data file. You can also save the auto-detected shape and use it later.
10. Click Save.
4-19
Chapter 4
Streams
• Tags
• Stream Type: The selected stream is displayed.
4. Click Next.
5. On the Source Type page, enter the following details:
• Connection: Select a GG Change Data.
• Table name: Enter a valid table name that includes the period (.) delimiter between
the catalog, schema, and table names. For example, [Link].table1
• Generate Full Records: Select this option to stream full data record (value of all
fields), irrespective of the database transactional changes to a single column, a
subset, or all the columns of a row.
– Database Connection: Select a GoldenGate sourced database connection.
– Enable Cache: Select this option to enable caching for GoldenGate Full Records,
to enhance its performance.
6. Click Next.
7. On the Shape screen, select one of the methods to define the shape:
• Infer Shape : Select this option to detect the shape automatically from the input data
stream.
• Select Existing Shape: Select one of the existing shapes from the drop-down list.
• Manual Shape : Select this option to manually infer the fields from a stream or file.
You can also update the datatype of the fields.
Note:
– To retrieve the entire JSON payload, add a new field with path $.
– To retrieve the content of the array, add a new field with path $
[arrayField].
Note:
The difference between a Kafka stream and a GoldenGate stream is that the pipeline
constructs, like the Query Group Table, understands the GoldenGate syntax and
associates it with the relevant GoldenGate fields.
4-20
Chapter 4
Streams
Note:
GGSA can read messages from Oracle Advanced Queue. This option is
available as a general JMS connector -
[Link].
• Jndi name: Enter the name of the Java interface that reads messages from topics,
distributed topics, queues and distributed queues
• Client ID: Enter the unique client ID to be used for a durable subscriber. If you do not
provide this value, subscriber ID is used as a clientID to create a durable subscriber.
• Message Selector : Set the message selector to filter messages. Message selectors
assign the work of filtering messages to the JMS provider rather than to the
application.
If your messaging application needs to filter the messages it receives, you can use a
JMS API message selector. A message selector is a String that contains an
expression. The syntax of the expression is based on a subset of the SQL92
conditional expression syntax. The message selector in the following example selects
any message that has a NewsType property that is set to the value Sports or Opinion:
4-21
Chapter 4
Streams
• Subscription ID: Enter the unique subscription ID for durable selector. This value is
essential for durable subscriber.
Note:
When you specify both clientID and subscriberID, you can have only one
running pipeline consuming that stream. If you need multiple subscribers/
pipelines, remove clientID or subscriberName from the stream or create
different streams (with different clientID and subscriberName) for multiple
pipelines.
• Data Format: Select the data format from the drop-down list. The supported formats
are: CSV, JSON, AVRO, MapMessage.
A MapMessage object is used to send a set of name-value pairs. The names are
String objects, and the values are primitive data types in the Java programming
language. The names must have a value that is not null, and not an empty string. The
entries can be accessed sequentially or randomly by name. The order of the entries is
undefined.
6. Click Next.
7. On the Data Format screen, enter the shape details for the stream, based on the data
format you have selected.
• For JSON:
– Allow Missing Column Names: Select this option to allow an input stream that
has a column undefined in the shape.
• For CSV:
– CSV Predefined Format: Select one of the predefined data format from the drop-
down list. For more information, see Predefined CSV Data Formats.
– First record as header: Select this option to use the first record as the header
row.
• For AVRO:
– Schema Namespace: Enter the schema name combined with the namespace, to
uniquely identify the schema within the store.
– Schema (optional): Upload a schema file to infer shape from.
• If you selected MapMessage as the data format, there are no specific attributes to be
set on this screen. The Data Format screen is skipped, and you are redirected to the
Shape screen.
8. Click Next.
9. On the Shape screen, select one of the methods to define the shape:
• Infer Shape : Select this option to detect the shape automatically from the input data
stream.
• Select Existing Shape: Select one of the existing shapes from the drop-down list.
• Manual Shape: Select this option to manually infer the fields from a stream or file. You
can also update the datatype of the fields.
4-22
Chapter 4
Streams
Note:
– To retrieve the entire JSON payload, add a new field with path $.
– To retrieve the content of the array, add a new field with path $
[arrayField].
4-23
Chapter 4
Streams
– Allow Missing Column Names: Select this option to allow an input stream that
has a column undefined in the shape.
• For CSV:
– CSV Predefined Format: Select one of the predefined data formats from the drop-
down list. For more information, see Predefined CSV Data Formats.
– First record as header: Select this option to use the first record as the header
row.
• For AVRO:
– Schema Namespace: Enter the schema name combined with the namespace, to
uniquely identify the schema within the store.
– Schema (optional): Upload a schema file to infer shape from.
8. Click Next.
9. On the Shape screen, select one of the methods to define the shape:
• Infer Shape : Select this option to detect the shape automatically from the input data
stream.
• Select Existing Shape: Select one of the existing shapes from the drop-down list.
• Manual Shape : Select this option to manually infer the fields from a stream or file.
You can also update the datatype of the fields.
Note:
– To retrieve the entire JSON payload, add a new field with path $.
– To retrieve the content of the array, add a new field with path $
[arrayField].
4-24
Chapter 4
Streams
2. On the Edit Stream screen, click the Edit link corresponding to the following sections, and
make the necessary changes.
• Source Details
• Source Type Parameters
• Data Type Parameters
• Source Shape
3. Click Save.
Deleting a Stream
To delete a stream:
1. On the Catalog page, hover the mouse over the stream that you want to delete.
2. Click the Deleteicon that appears to your right side on the screen.
3. On the Delete Confirmation screen, click Delete.
4-25
Chapter 4
Streams
Note:
The input timestamp is truncated to millisecond precision.
4-26
Chapter 4
References
4.3 References
Create References
Manage References
4-27
Chapter 4
References
• Name: Enter a unique name for the reference. This is a mandatory field.
• Display Name: Enter a display name for the reference. If left blank, the Name field
value is copied.
• Description
• Tags
• Reference Type: The selected reference is displayed.
4. Click Next.
5. On the Source Details page, provide the following details:
• Connection: Select a coherence connection for the reference
.
• Cache name: Enter the name of the coherence cache to enable caching. Caching is
supported only for single equality join condition.
• Data Format: Select POJO or Map as the data format for the reference.
6. Click Next.
7. On the Shape screen, select one of the methods to define the shape :
• If you selected Map as the data format, you have the following options to define the
shape:
Select Existing Shape: Select an existing shape that you want to use for the
reference.
• Manual Shape: Select this option if you want to define your own shape.
8. Click Save.
For information on data mapping in the two coherence reference types, see:
• Data Mapping in Coherence Reference Map Type
• Data Mapping in Coherence Reference POJO Type
4-28
Chapter 4
References
• Enable Caching: Select this option to enable caching. Ignite cache is the default
cache for cache-enabled references. So before deploying a pipeline using cache-
enabled references, you have to start the cache cluster from the System Settings tab.
Note:
The Enable Caching option is not supported in the GGSA marketplace
instance.
• Caching Scheme: Select the cache type from the drop-down list.
In a Partitioned Cache the data is partitioned among all the machines of the cluster.
In a Replicated Cache the data is fully replicated to every member of the cluster. Use
Replicated Cache when the number of cache entries are relatively low and do not
need to be updated often.
• Expiry Delay: The duration delay from last update that the entries will be kept by the
cache before being marked as expired. Any attempt to read an expired entry will result
in a reloading of the entry from the configured cache store. This field is enabled only
when caching is enabled.
5. Provide details for the following fields on the Shape page and click Save:
• Shape Name: Select a shape that you want to use for the reference
When the datatype of the table data is not supported, the table columns do not have auto
generated datatype. Only the following datatypes are supported:
• numeric
• interval day to second
• text
• interval year to month
• timestamp (without timezone)
• date time (without timezone)
Note:
The date column cannot be mapped to timestamp. This is a limitation in the
current release.
4-29
Chapter 4
References
• Description
• Tags
• Reference Type: The selected reference is displayed.
4. Click Next.
5. On the Source Details page, provide the following details:
• Connection: Select an ignite cache connection for the reference.
• Cache name: Enter the name of the ignite cache to enable caching. Caching is
supported only for single equality join condition.
• Data Format: Select the data format from the drop-down list.
6. Click Next.
7. On the Shape screen, select one of the methods to define the shape:
• Infer Shape: Select this option to detect the shape automatically from the input data
stream.
– From Stream: Select this option to infer shape from a stream.
– From File: Select this option to infer the shape from a JSON schema file, or a
JSON or CSV data file. You can also save the auto-detected shape and use it
later.
• Select Existing Shape: Select one of the existing shapes from the drop-down list.
• Manual Shape: Select this option to manually edit the shape. You can also update the
datatype of the fields.
8. Click Save.
Deleting a Reference
To delete a reference:
1. Go to the Catalog page and hover the mouse over the reference that you want to delete.
2. Click the delete icon that appears to your right side on the screen.
3. On the Delete Confirmation screen, click Delete.
4-30
Chapter 4
References
<caching-schemes>
.............
<proxy-scheme>
<service-name>ExtendTcpProxyService</service-name>
<acceptor-config>
<tcp-acceptor>
<local-address>
<address>{ADDRESS}</address>
<port>{PORT}</port>
</local-address>
</tcp-acceptor>
</acceptor-config>
<autostart>true</autostart>
</proxy-scheme>
</caching-schemes>
4-31
Chapter 4
References
[Link]("strValue", "Test");
[Link]("intervalValue", "+000000002 03:04:11.330000000");
[Link]("orderTag", big10);
[Link](big10,order1);
Note:
When you upload the POJO in a jar, you must ensure that the fully-qualified class
name of POJO matches exactly in cache and the custom POJO jar.
For example, if you have loaded [Link] objects in cache, custom
JAR should have CustomPOJO class inside the [Link] package.
If there any mismatch in the class name, then querying the cache for certain type of
objects, does not return a result and you will not see any data in the live output table.
4-32
Chapter 4
References
• [Link] (Float)
• boolean (Boolean)
• [Link] (Boolean)
• [Link]
("externalcachepojo");
Object> cache) {
For this example coherence reference should be created with cache name as
externalcachepojo and can join with stream with orderid.
4-33
Chapter 4
Targets
return orderId;
}
public String getOrderDesc() {
return orderDesc;
}
public boolean equals(Object object) {
if (this == object) return true;
if (object == null || getClass() != [Link]()) return false;
if () return false;
OrderPOJO that = (OrderPOJO) object;
return [Link](orderId, [Link]) &&
[Link](orderDesc, [Link]);
}
public int hashCode() {
return [Link]([Link](), orderId, orderDesc);
}
}
Note:
Ensure that the POJO class does not have a GGSA coherence target as a
constructor, because it can instantiate the POJO class using default constructor, and
then access the setXXX and getXXX, and isXXX methods.
4.4 Targets
Create Targets
Manage Targets
4-34
Chapter 4
Targets
• MongoDB Target
• NFS Target
• Notification Target
• Object Storage Target
• OSS Target
• REST Target
4-35
Chapter 4
Targets
* Avro Codec: Select a compression codec from the drop-down list. This option
is enabled if you have selected the file format as AVRO or AVRO Object
Container Format.
• For PARQUET:
– PARQUET Compression: Select a compression codec from the drop-down list.
• For ORC:
– ORC Compression: Select a compression codec from the drop-down list.
8. Click Next.
9. On the Shape screen, select one of the methods to define the shape:
• Select Existing Shape: Select one of the existing shapes from the drop-down list.
• Manual Shape: Select this option to manually infer the fields from a stream or file. You
can also add to, or remove fields from, an existing shape. Enter the following details:
– Shape Name: Enter a name for the shape.
– Clear Fields: Click to delete all the fields in the shape.
– Field Name: Add the necessary fields.
– Field Type: Select the field data type from the drop-down list.
10. Click Save.
4-36
Chapter 4
Targets
• NFS Path: Enter the local file or NFS path where the files are written first and then
uploaded to HDFS.
• Storage Format: Select a storage format from the drop-down list.
6. Click Next.
7. On the Data Format screen, enter the shape details, based on the storage format you
have selected.
• For FILE:
– File Format: Select a file format from the drop-down list.
* JSON Delimiter: Enter the JSON delimiter if you have selected the JSON file
format.
* Avro Codec: Select a compression codec from the drop-down list. This option
is enabled if you have selected the file format as AVRO or AVRO Object
Container Format.
• For PARQUET:
– PARQUET Compression: Select a compression codec from the drop-down list.
• For ORC:
– ORC Compression: Select a compression codec from the drop-down list.
8. Click Next.
9. On the Shape screen, select one of the methods to define the shape:
• Select Existing Shape: Select one of the existing shapes from the drop-down list.
• Manual Shape: Select this option to manually infer the fields from a stream or file. You
can also add to, or remove fields from, an existing shape. Enter the following details:
– Shape Name: Enter a name for the shape.
– Clear Fields: Click to delete all the fields in the shape.
– Field Name: Add the necessary fields.
– Field Type: Select the field data type from the drop-down list.
10. Click Save.
4-37
Chapter 4
Targets
Note:
* To retrieve the entire JSON payload, add a new field with path $.
* To retrieve the content of the array, add a new field with path $
[arrayField].
4-38
Chapter 4
Targets
2. Hover the mouse over Target and select Database from the submenu.
3. On the Type Properties screen, enter the following details:
• Name: Enter a unique name for the target. This is a mandatory field.
• Display Name: Enter a display name for the target. If left blank, the Name field value
is copied.
• Description
• Tags
• Target Type: The selected target is displayed.
4. Click Next.
5. On the Target Details screen, enter the following details:
• Connection: Select a database connection from the drop-down list.
6. Click Next.
7. On the Shape screen, enter the following details:
• Table Name: Select a database table from the drop-down list.
8. Click Save.
4-39
Chapter 4
Targets
Note:
* Any update to the value will result in new entry rather than updating
previous value.
* If a key has a null value, ElasticSearch will autogenerate the key. In
the example above, ID is 2, because serial is selected as the key
field. If record has null in serial field:
{"address":"Mumbai","serial":null","clientName":"Joe"}, then
ID will be autogenerated by Elasticsearch.
* Index is json_data which is provided in previous step, ID will be value
of each record.
4-40
Chapter 4
Targets
• Select Existing Shape: Select one of the existing shapes from the drop-down list.
• Manual Shape: Select this option to manually define a shape. You can also add to, or
remove fields from, an existing shape. Enter the following details:
– Shape Name: Enter a name for the shape.
– Clear Fields: Click to delete all the fields in the shape.
– Key: Selecting atleast one field in the HBase table as a primary key. A primary key
is mandatory.
– Field Name: Add the necessary fields.
– Field Type: Select the field data type from the drop-down list.
8. Click Save.
4-41
Chapter 4
Targets
* JSON Delimiter: Enter the JSON delimiter if you have selected the JSON file
format.
* Avro Codec: Select a compression codec from the drop-down list. This option
is enabled if you have selected the file format as AVRO or AVRO Object
Container Format.
• For PARQUET:
– PARQUET Compression: Select a compression codec from the drop-down list.
• For ORC:
– ORC Compression: Select a compression codec from the drop-down list.
8. Click Next.
9. On the Shape screen, select one of the methods to define the shape:
• Select Existing Shape: Select one of the existing shapes from the drop-down list.
• Manual Shape: Select this option to manually infer the fields from a stream or file. You
can also add to, or remove fields from, an existing shape. Enter the following details:
– Shape Name: Enter a name for the shape.
– Clear Fields: Click to delete all the fields in the shape.
– Field Name: Add the necessary fields.
– Field Type: Select the field data type from the drop-down list.
10. Click Save.
4-42
Chapter 4
Targets
• File Roll Interval: Enter the roll-over interval to write a new file. The interval can be in
10ms, 10s, 10m, 1hr formats.
• File Roll Max Size: Enter the roll-over file size to create a new file. The size can be in
1000, 10k, 10m, 1g formats.
6. Click Next.
7. On the Shape screen, select one of the methods to define the shape:
• Select Existing Shape: Select one of the existing shapes from the drop-down list.
• Manual Shape: Select this option to manually infer the fields from a stream or file. You
can also add to, or remove fields from, an existing shape. Enter the following details:
– Shape Name: Enter a name for the shape.
– Key: Select key fields, based on which the data is partitioned.
– Clear Fields: Click to delete all the fields in the shape.
– Field Name: Add the necessary fields.
– Field Type: Select the field data type from the drop-down list.
8. Click Save.
4-43
Chapter 4
Targets
• Backup: Enter the number of backup nodes. This field is not applicable for an
embedded cluster.
• Update Cache Entry: Select this option to update a particular key value with new
data. This option is enabled by default.
• Data Format: Select a data format from the drop-down list.
6. Click Next.
7. On the Data Format screen, enter the shape details for the stream, based on the data
format you have selected.
• For JSON:
– Create nested json object: Select this option to create a nested JSON object for
the target.
8. Click Next.
9. On the Shape screen, enter the following details:
• For JSON:
– Shape Name: Enter a name for the shape.
– Clear Fields: Click to clear all the fields from the existing shape.
– Key: Select a key from the input data to store record.
You can select multiple fields as key. Key selection is mandatory.
– Field Name: Add the necessary fields.
– Field Path: Enter the field path.
– Field Type: Select the field data type from the drop-down list.
10. Click Save.
Note:
You cannot edit an Ignite target once created. This restriction avoids cache data
corruption because only one target from the GGSA platform is allowed to write to only
one cache in the ignite server.
4-44
Chapter 4
Targets
Note:
* To retrieve the entire JSON payload, add a new field with path $.
* To retrieve the content of the array, add a new field with path $
[arrayField].
4-45
Chapter 4
Targets
* Field Type: Select the field data type from the drop-down list.
– For CSV, AVRO, and MapMessage:
* Shape Name: Enter a name for the shape.
* Clear Fields: Click to delete all the fields from the shape.
* Field Name: Add the necessary fields.
* Field Type: Select the field data type from the drop-down list.
10. Click Save.
4-46
Chapter 4
Targets
• Select Existing Shape: Select one of the existing shapes from the drop-down list.
• Manual Shape: Select this option to manually infer the fields from a stream or file. You
can also add to, or remove fields from, an existing shape. Enter the following details:
– For JSON:
* Shape Name: Enter a name for the shape.
* Clear Fields: Click to remove all the fields from the shape.
* Key: Select key fields, based on which data is partitioned. For example,
records containing the same values for the selected key fields will all be stored
in the same Kafka partition.
You can select multiple fields as key. Key selection is not mandatory.
* Field Name: Add the necessary fields.
* Field Path: Enter the field path.
Note:
* To retrieve the entire JSON payload, add a new field with path $.
* To retrieve the content of the array, add a new field with path $
[arrayField].
4-47
Chapter 4
Targets
4. Click Next.
5. On the Target Details screen, enter the following details:
• Connection: Select a MongoDB connection from the drop-down list.
• Database: Enter the name of the database to be used for the target.
• Collection: Enter the name of the collection to insert documents.
6. Click Next.
7. On the Shape screen, select one of the methods to define the shape:
• Select Existing Shape: Select one of the existing shapes from the drop-down list.
• Manual Shape: Select this option to manually define a shape. You can also add to, or
remove fields from, an existing shape. Enter the following details:
– Shape Name: Enter a name for the shape.
– Clear Fields: Click to delete all the fields in the shape.
– Key: Select none, or one or more fields as key. This key will be the ID field.
– Field Name: Add the necessary fields.
– Field Type: Select the field data type from the drop-down list.
8. Click Save.
4-48
Chapter 4
Targets
7. On the Storage Format screen, enter the shape details, based on the storage format you
have selected.
• For FILE:
– File Format: Select a file format from the drop-down list.
– Avro Codec: Select a compression codec from the drop-down list. This option is
enabled if you have selected the file format as AVRO or AVRO Object Container
Format.
• For PARQUET:
– PARQUET Compression: Select a compression codec from the drop-down list.
• For ORC:
– ORC Compression: Select a compression codec from the drop-down list.
8. Click Next.
9. On the Shape screen, select one of the methods to define the shape:
• Select Existing Shape: Select one of the existing shapes from the drop-down list.
• Manual Shape: Select this option to manually infer the fields from a stream or file. You
can also add to, or remove fields from, an existing shape. Enter the following details:
– Shape Name: Enter a name for the shape.
– Clear Fields: Click to delete all the fields in the shape.
– Field Name: Add the necessary fields.
– Field Type: Select the field data type from the drop-down list.
10. Click Save.
Note:
In case of JSON, AVRO, AVRO OCF schema files would be written under NFS
Path/SCHEMA folder.
4-49
Chapter 4
Targets
Note:
ONS connection type is no longer supported in GGSA. Recreate older Notification
type targets in the pipeline, using an OCI connection.
4-50
Chapter 4
Targets
• NFS Path: Enter the local file or NFS path where the files are written first and then
uploaded to the Object Storage.
• Storage Format: Select a storage format from the drop-down list.
6. Click Next.
7. On the Data Format screen, enter the shape details, based on the storage format you
have selected.
• For FILE:
– File Format: Select a file format from the drop-down list.
* JSON Delimiter: Enter the JSON delimiter if you have selected the JSON file
format.
* Avro Codec: Select a compression codec from the drop-down list. This option
is enabled if you have selected the file format as AVRO or AVRO Object
Container Format.
• For PARQUET:
– PARQUET Compression: Select a compression codec from the drop-down list.
• For ORC:
– ORC Compression: Select a compression codec from the drop-down list.
8. Click Next.
9. On the Shape screen, select one of the methods to define the shape:
• Select Existing Shape: Select one of the existing shapes from the drop-down list.
• Manual Shape: Select this option to manually infer the fields from a stream or file. You
can also add to, or remove fields from, an existing shape. Enter the following details:
– Shape Name: Enter a name for the shape.
– Clear Fields: Click to delete all the fields in the shape.
– Field Name: Add the necessary fields.
– Field Type: Select the field data type from the drop-down list.
10. Click Save.
4-51
Chapter 4
Targets
4. Click Next.
5. On the Target Details screen, enter the following details:
• Connection: Select a Kafka connection.
• Topic name: Enter a name for the kafka topic.
• Data Format: Select a data format from the drop-down list.
6. Click Next.
7. On the Data Format screen, enter the shape details, based on the data format you have
selected.
• For JSON:
– Create nested json object: Select this option to create a nested JSON object for
the target.
• For CSV:
– CSV Predefined Format: Select one of the predefined data formats from the drop-
down list. For more information, see Predefined CSV Data Formats.
– First record as header: Select this option to use the first record as the header
row.
8. Click Next.
9. On the Shape screen, select one of the methods to define the shape:
• Infer Shape: Select this option to detect the shape automatically from the input data
stream.
• Select Existing Shape: Select one of the existing shapes from the drop-down list.
• Manual Shape: Select this option to manually infer the fields from a stream or file. You
can also add to, or remove fields from, an existing shape. Enter the following details:
– For JSON:
* Shape Name: Enter a name for the shape.
* Clear Fields: Click to clear all the fields from the existing shape.
* Key: Select key fields, based on which data is partitioned. For example,
records containing the same values for the selected key fields will all be stored
in the same Kafka partition.
You can select multiple fields as key. Key selection is not mandatory.
* Field Name: Add the necessary fields.
* Field Path: Enter the field path.
Note:
* To retrieve the entire JSON payload, add a new field with path $.
* To retrieve the content of the array, add a new field with path $
[arrayField].
4-52
Chapter 4
Targets
Note:
The Trust Store File and Trust Store Password options allow the use of
untrusted certificates for REST connections, resulting in an insecure
connection.
4-53
Chapter 4
Targets
• Batch processing: Select this option to process batch events as a single request.
Enable this option for high throughput pipelines. For example,
Eg: [{"address":
{ "street" : xxxxxxx }
},{"address":
{ "street" : xxxxxxa }]
• HTTP Method: Select this option for the REST target to send requests to REST end-
point, using Http POST and PUT methods. Default is set to POST.
• Data Format: Select a data format from the drop-down list.
6. Click Next.
7. On the Data Format screen, enter the shape details, based on the data format you have
selected.
• For JSON:
– Create nested json object: Select this option to create a nested JSON object for
the target. For example, if the target shape is defined as
field:attribute_street, field_path:address/street.
}
{"address":
{ "street" : xxxxxxx }
• For CSV:
– CSV Predefined Format: Select one of the predefined data formats from the drop-
down list. For more information, see Predefined CSV Data Formats.
– First record as header: Select this option to use the first record as the header
row.
8. Click Next.
9. On the Shape screen, select one of the methods to define the shape:
• Infer Shape: Select this option to detect the shape automatically from the input data
stream.
• Select Existing Shape: Select one of the existing shapes from the drop-down list.
• Manual Shape: Select this option to manually infer the fields from a stream or file. You
can also add to, or remove fields from, an existing shape. Enter the following details:
– For JSON:
* Shape Name: Enter a name for the shape.
* Clear Fields: Click to remove all the fields from the shape.
* Key: Select key fields, based on which data is partitioned. For example,
records containing the same values for the selected key fields will all be stored
in the same Kafka partition.
You can select multiple fields as key. Key selection is not mandatory.
* Field Name: Add the necessary fields.
4-54
Chapter 4
Targets
Note:
* To retrieve the entire JSON payload, add a new field with path $.
* To retrieve the content of the array, add a new field with path $
[arrayField].
4-55
Chapter 4
Pipelines
}
public String setOrderDesc(String str2) {
[Link]=str2;
}
Note:
Ensure that the POJO class does not have a GGSA coherence target as a
constructor, because it can instantiate the POJO class using default constructor, and
then access the setXXX and getXXX, and isXXX methods.
Class Abc {
4.5 Pipelines
Create a Pipeline
Manage Pipelines
4-56
Chapter 4
Pipelines
4-57
Chapter 4
Pipelines
Note:
Make sure to allot more memory to executors in the scenarios where you have
large windows.
To export a pipeline:
1. On the Catalog page, hover the mouse over, or select the pipeline that you want to export
to another GGSA instance.
2. Click the Export option that appears to your right side on the screen.
3. The selected pipeline and its dependent artifacts are exported as a JSON zip file, to your
computer's default Downloads folder.
To import a pipeline:
1. Go to the GGSA instance to which you want to import the exported metadata.
2. On the Catalog page, click Import.
3. In the Import dialog box, click Select, to locate and select the exported zip file on your
computer.
4-58
Chapter 4
Pipelines
4. On the Import Resources tab, you can select an existing connection from the catalog, or
use the imported connection.
5. Click Import.
The imported pipeline and its dependent artifacts are available on the Catalog page.
Note:
• Each pipeline should have a unique name. If you are importing an updated
version of a pipeline, you can retain the same name. If you are importing a new
pipeline and if a pipeline with the same name already exists in the catalog,
change the name of the pipeline that you are importing.
• If you have already exported a pipeline with the same name, update the pipeline
name as below:
1. Create a directory exportUpdate.
2. Copy the exported zip, say [Link], to the folder
exportUpdate.
3. Unzip the file [Link].
4. Open the json file in edit mode.
5. Search for pipeline/ artifact name in the json file. For example, if Nano
pipeline was the name given to the pipeline, update it to Nano pipeline
updated.
6. Update the json file in [Link].
7. Import this zip.
8. The pipeline is automatically assigned a name, using the display name.
9. The draft pipeline and publish pipeline topic are created as below:
a. sx_Nanopipelineupdated_Nano_Stream_draft
b. sx_Nanopipelineupdated_Nano_Stream_public
4-59
Chapter 4
Pipelines
Hide/Unhide Columns
In the live output table, right-click a column and click Hide to hide that column from the output.
This option only hides the columns from the UI and does not remove them from the output. To
unhide the hidden columns, click Columns and then click the eye icon to make the columns
visible in the output.
Add a Timestamp
Include timestamp in the live output table by clicking the clock icon in the output table.
4-60
Chapter 4
Pipelines
Click the Show Topology icon at the top-right corner of the editor to open the topology viewer.
By default, the topology of the entity from which you launch the Topology Viewer is displayed.
The context of this topology is Immediate Family, which indicates that only the immediate
dependencies and connections between the entity and other entities are shown. You can
switch the context of the topology to display the full topology of the entity from which you have
launched the Topology Viewer. The topology in an Extended Family context displays all the
dependencies and connections in the topology in a hierarchical manner.
Note:
The entity for which the topology is shown has a grey box surrounding it in the
Topology Viewer.
Immediate Family
Immediate Family context displays the dependencies between the selected entity and its child
or parent.
The following figure illustrates how a topology looks in the Immediate Family.
Extended Family
Extended Family context displays the dependencies between the entities in a full context, that
is if an entity has a child entity and a parent entity, and the parent entity has other
dependencies, all the dependencies are shown in the Full context.
The following figure illustrates how a topology looks in the Extended Family.
4-61
Chapter 4
GoldenGate Change Stream
Note:
The GoldenGate username and password of the deployment should be of the
user with access to create a new distribution path from the Goldengate instance.
4-62
Chapter 4
GoldenGate Change Stream
5. Click Next.
6. On the GG Change Data Details page, enter the following details:
• GG Extracts: Select a GG stream from the drop-down list.
• Target Trail: Enter a two character name for the Goldengate trail file.
• Kafka Connection: Select a Kafka connection from the drop-down list.
• GG Change Data name: Enter a name for the goldengate stream (maximum 8
characters). This name will be used for the replicat process that puts the change data
from trail file to Kafka topics.
7. Click Save.
Note:
The following template parameter files for the replicat process are located at osa-
base/etc/:
• [Link]
• [Link]
• custom_kafka_producer.[Link]
You can modify these template files to customize the replicat process before
proceeding to the next step.
4-63
Chapter 4
GoldenGate Change Stream
Note:
When you start a GG Change Data replicat process, it creates kafka topics, and
starts pushing changed data to the new topics. For example, if there are 10 tables in
the extract process that you chose while creating the GG Change Data, 10 new
topics will be created.
The names of the topics created are in the following format:
GGChangeDataName_fullyQualifiedTableName
You can use these topics to create a new stream (with Goldengate as stream type),
and in pipelines, similar to using a Kafka stream.
The trail files, after being processed completely by the replicat process, and after one hour of
inactivity, will be purged. The files will be checked for purging every 10 minutes.
You can modify the above rule, in an OCI GGSA VM, following the steps below:
1. Stop the manager process by running the command sudo systemctl stop ggbd-
mgr.
2. Modify the rules in the file /u01/app/ggbd/OGG_BigData_Linux_x64_19.[Link]/
dirprm/[Link].
3. Start the manager process by running the command sudo systemctl start ggbd-
mgr.
For more information on rules about purging the trail files, see PURGEOLDEXTRACTS for
Manager.
4-64
Chapter 4
Embedded Ignite Cache
Note:
Full Records option is not supported in the GGSA marketplace instance.
Note:
Ignite Caching is not supported in the GGSA marketplace instance.
Note:
NFS mounted path is the preferred persistence store path, to enable cache
rehydration, on restart of a cluster or a node.
4-65
Chapter 4
Embedded Ignite Cache
Note:
When you start a cache cluster for the first time, an embedded Ignite Connection is
created with connection type as Ignite Cache and Connection Name as Embedded.
This connection is not editable connection, and connection details are not shown, but
you can select this connection while creating an Ignite Target and an Ignite
Reference.
4-66
Chapter 4
Ignite Cluster on OCI GGSA
The Cluster Status changes to In Progress. You will see a status message to refresh the
page. Close the System Settings page and reopen it. The Cluster Status changes to
Running.
4-67
Chapter 4
GGBD Cluster on OCI GGSA
Note:
• The Delete Storage option is available only while the Ignite Cluster is running.
• If you leave the box unselected, the cached values are retained for the next
restart.
4-68
5
Transform
Correlating an event with other sources requires the join condition to be based on a common
key. In the above example, the SensorId from the stream cand be used to correlate with
SensorKey in the database table. The following query illustrates the above data enrichment
scenario producing sensor details for all sensors whose temperature exceeds their pre-defined
threshold.
Queries like above and more complex queries can be automatically generated by configuring
sources and filter sections of the query stage.
Stream-to-stream Correlation
• A Stream is an unbounded sequence of events. To correlate a Stream with another
Stream, first convert both streams to a Relation or a bounded sequence of events, by
applying window functions.
• After applying window functions on both streams, define a correlation condition that
evaluates to true or false.
The output from stream-to-stream correlation is a subset of the Cartesian product of tuples
from both windows, where the correlation condition is true.
5-1
Chapter 5
Applying Window Functions to a Stream
Stream-to-Cache Correlation
• Convert the Stream to a bounded sequence of events, by applying a window function.
• After applying the window function on the stream, define a correlation condition that
evaluates to true of false.
The output from Stream-to-Cache is the Cartesian product of tuples from window and the
cache, where the correlation condition is true. Currently, OSA supports only Coherence cache.
In the above example, the data is retained for 5 minutes but the query is evaluated every 30
seconds.
5-2
Chapter 5
Applying Window Functions to a Stream
Note:
There will be an output only if the current results of the query is different from the
previous results. This avoids sending duplicates to downstream applications.
If you set the slide to same as range, it will create a tumbling window instead of a sliding
window. For example,
will only retain 5 minutes of data and the query will only execute every 5 minutes.
Note:
Use the tumbling window to batch output results before sending to downstream
systems. For example, you may want to create a file on object store only after you
have accumulated at least 10000 events. This would avoid the small-file problem in
Big-Data systems. Similarly, you may want to avoid multiple writes to a database
system and instead perform a single write, after sufficient events have been
accumulated.
[range 1 minutes ]
In the example above, the default slide value which is same as spark streaming batch interval
is used.
Note:
If you do not specify a slide value, it will take the default slide, which is same as the
Spark batch interval.
[rows 10 slide 1]
5-3
Chapter 5
Applying Window Functions to a Stream
Maximum window size is 10 events, but Slide of 1 implies the query is executed on the arrival
of every new event.
[rows 10]
Last 10 events is used to evaluate the query. Default slide value is used.
[ CurrentYear ]
Data is retained until end of the current year. Default slide value is used.
CurrentMonth
• Applicable on: Query Stage, Detect Duplicates Pattern, Eliminate Duplicate Pattern
The CQL example is as follows:
[ CurrentMonth ]
Data is retained until the end of the current month. Default slide value is used.
CurrentDay
• Applicable on: Query Stage, Detect Duplicates Pattern, Eliminate Duplicate Pattern
The CQL example is as follows:
[ CurrentDay ]
CurrentHour
• Supported types of shape fields: timestamp, int, bigint
• Applicable on: Query Stage, Detect Duplicates Pattern, Eliminate Duplicate Pattern
The CQL example is as follows:
[ CurrentHour ]
Data is retained until the end of the current hour. Default slide value is used.
5-4
Chapter 5
Applying Window Functions to a Stream
will only retain events from the last 2 days, based on the timestamp value in
EventCaptureTime field.
Last 10 events for each partition value. For example [partition by stockSymbol rows 10] will use
last 10 quotes for ORCL, last 10 quotes for AMZN, etc.
Query is evaluated on the arrival of new events and not on time ticks.
Default slide value is used.
5.3.8 Applying a Row Window with Partition with Range without Slide
• Shape fields of Partition by: MultiSelect
• Rows value: Integer
• Range value: Integer
• Range unit: nanoseconds, microseconds, milliseconds, seconds, minutes, hours
• Applicable on: Query Stage
The CQL example is as follows:
Events may be evicted from the window even when it is not full with all 10 rows, but 15
seconds have elapsed since the event arrived.
5-5
Chapter 5
Adding Stages to a Pipeline
5.3.9 Applying a Row Window with Partition with Slide and Range
• Shape fields of Partition by: MultiSelect
• Rows value: Integer
• Range value: Integer
• Range unit: nanoseconds, microseconds, milliseconds, seconds, minutes, hours
• Slide Value: Integer
• Slide unit: nanoseconds, microseconds, milliseconds, seconds, minutes, hours
• Applicable on: Query Stage
The CQL example is as follows:
5-6
Chapter 5
Adding Stages to a Pipeline
Note:
IN operator is available as an operator in the drop-down list. This operator is not
supported for Interval, Interval YM, Timestamp, and SDO Geometry datatypes.
You can use the IN filter to refer to a column in a database table. When you
change the database column values at runtime, the pipeline picks up the latest
values from the DB column, without republishing the pipeline.
5-7
Chapter 5
Adding Stages to a Pipeline
5-8
Chapter 5
Adding Stages to a Pipeline
11. On the Visualizations tab, click Add a Visualization and add the required type of
visualization. See Adding Chart Visualizations.
5-9
Chapter 5
Adding Stages to a Pipeline
6. Select a suitable condition in the IF statement, THEN statement, and click Add Action to
add actions within the business rules.
Actions can also be expressions. For example, SET Revenue TO =-Revenue, will convert
the current value of Revenue to a negative number.
Expressions must always start with a '=' sign. For a constant text value, just type in the
text. For example, SET CustomerType TO GOLD.
The rules are applied to the incoming events one by one and actions are triggered if the
conditions are met.
5-10
Chapter 5
Applying Functions to Create a New Column
Note:
Currently, you can use expressions only within a query stage.
Using Functions
You can select a CQL Function from the list of available functions and select the input
parameters. Make sure to begin the expression with ”=”. Click Apply to apply the function to
the streaming data.
Example expression using functions:
=float((CanceledOrdersFloat/NewOrdersFloat) * 100.0)
5-11
Chapter 5
Applying Functions to Create a New Column
You can see custom functions in the list of available functions when you add/import a custom
jar in your pipeline.
For a list of supported functions, see #unique_202 .
5-12
Chapter 5
Applying Functions to Create a New Column
[Link] BesselI0
Returns the modified Bessel function of order 0 of the input argument.
The input arguments can be one of the following data types: double, float.
Returned value type will be double.
Function Result
besselIO(65) 8.403039845625433E26
besselIO(3125.2) 1.07389541368045088E17
[Link] BesselIO_exp
Returns the exponentially scaled modified Bessel function of order 0 of the double argument as
a double.
The input argument can be one of the following data type: double, integer, float. Returned
value type will be double.
Function Result
besselIO_exp(1451.44) 8.113723742037748E23
[Link] BesselI1(value1)
Function returns the modified Bessel function of order 1 of the double argument.
The input arguments can be one of the following data types: double, integer, float.
The returned value type will be double.
Function Result
besselI1(432.98) 2.1043808863643512E186
besselI1(31) 2.055972795294565E12
[Link] BesselI1_exp(value1)
Function returns the exponentially scaled modified Bessel function of order 1, of the input
argument.
The input arguments can be one of the following types: double, integer, float.
Returned value type will be double.
5-13
Chapter 5
Applying Functions to Create a New Column
Function Result
besselI1_exp(99) 0.03994284829937756
[Link] BesselK0_exp(value1)
Function returns the exponentially scaled modified Bessel function of the third kind of order 0.
Input value can be one of the following types: double, integer, float.
Returned value type will be double.
Function Result
Besselk0_exp(3.6) 0.6404559726736455
[Link] BesselIK1_exp(value1)
Function returns the exponentially scaled modified Bessel function of the third kind of order 1.
Input value can be one of the following types: double, integer, float.
Returned value type will be double.
Function Result
BesselIK1_exp(72) 0.14847048263652857
BesselIK1_exp(3.6) 0.7244606719817783
Function Result
BesselY(30,2.2) -1.6816755062290252E29
5-14
Chapter 5
Applying Functions to Create a New Column
Function Result
besselJ(4,3.3) 0.1742753869717833
[Link] BesselK(value1,value2)
Function returns the modified Bessel function of the third kind of order n of the input argument.
Value1 can be one of the following types: integer.
Value 2 can be one of the following types: double, integer, float.
Returned value type will be double.
Function Result
BesselK(30,2) 4.271125754887687E30
[Link] bigdecimal(value1)
Converts the input argument value to big decimal. The input argument can be one of the
following data types: big integer, number, double, integer, text, float. Returned value type will
be number.
Function Result
bigdecimal(60) 6E+1
bigdecimal(32) 32
[Link] boolean(value1)
Converts the input argument value to logical. The input argument can be one of the following
data type: big integer or integer. Returned value type will be Boolean.
5-15
Chapter 5
Applying Functions to Create a New Column
Examples
Function Result
boolean(5) TRUE
boolean(0) FALSE
boolean(NULL) TRUE
boolean() TRUE
boolean(-5) TRUE
[Link] double(value1)
Converts the input argument value to double. The input argument can be one of the following
data types: integer, big integer, double, text or float. Returned value type will be double.
Examples
Function Result
double(3.1406) 3.1405999660491943
double(1234.56) 1234.56005859375
[Link] float(value1)
Converts the input argument value to float. The input argument can be one of the following
data types: integer, big integer, double, text or float. Returned value will be a single-precision
floating-point number.
Examples
Function Result
float(1.67898989395) 1.6789899
float(1.796709289) 1.7967093
float(12.60508090750) 12.605081
[Link] int(value1)
Converts the input argument value to integer. The input argument can be one of the following
types: integer, text. Returned value type will be integer.
Function Result
int(50/3) 16
[Link] long()
Converts the input argument value to long. The input argument can be one of the following
types: big integer, integer, text, float, timestamp. Returned value type will be big integer.
5-16
Chapter 5
Applying Functions to Create a New Column
Function Result
long(5039505078907524) 5039505078907524
long(22) 22
Function Result
string(transaction_time,"hh-mm-ss") , 12-23-04
where transaction_time is 12/19/2016 12:23:04
string(transaction_time,"M-DD-YY") , 12-19-16
where transaction_time is 12/19/2016 12:23:04
5-17
Chapter 5
Applying Functions to Create a New Column
5-18
Chapter 5
Applying Functions to Create a New Column
[Link] Day(date)
day(date) function takes as an argument any one of the following data types: time interval or
timestamp. The returned value represents the day in the timestamp represented by this date
object. Returns a big integer indicating the day represented by this date.
Examples
Function Result
day(transaction-time), where 19
transaction_time is 12/19/2016 12:22:48
[Link] eventtimestamp(value1)
Event timestamp from stream.
Returned value will be of type timestamp.
Function Result
eventtimestamp() 4/4/2019 16:40:57
,
[Link] hour(date)
hour(date) function takes as an argument any one of the following data types: time interval or
timestamp. The returned value represents the hour in the time represented by this date object.
Returns a big integer indicating the hour of the time represented by this date.
Examples
Function Result
hour(12/06/17 09:15:22 AM) 09
hour(2015:07:21 12:45:35 PM) 12
[Link] minute(date)
minute(date) function takes as an argument any one of the following data types: time interval
or timestamp. The returned value represents the minutes in the time represented by this date
object. Returns a big integer indicating the minutes of the time represented by this date.
5-19
Chapter 5
Applying Functions to Create a New Column
Examples
Function Result
minute(12/06/17 09:15:22 AM) 15
minute(2015:07:21 12:45:35 PM) 45
[Link] month(date)
month(date) function takes as an argument any one of the following data types: time interval
or timestamp. The returned value represents the month of the year that contains or begins with
the instant in time represented by this date object. Returns a big integer indicating the month of
the year represented by this date.
Examples
Function Result
month(12/06/17 09:15:22 AM) 12
month(2017:09:23 11:20:25 AM) 9
[Link] nanosecond(value1)
Extracts and returns the current fractional part of second from date.
Value 1 can be one of the following types: timestamp.
Returned value will be of type big integer.
Function Result
nanosecond(transaction_time), 12/19/2016 719978080
12:22:57
[Link] systemtimestamp(value1)
Returns the current system time.
Returned value will be of type timestamp.
Function Result
systemtimestamp() 4/4/2019 17:06:14
,
5-20
Chapter 5
Applying Functions to Create a New Column
Function Result
timeformat(transaction_time, "M-dd-yy"), 12-19-16
where calc is 12/19/2016 12:22:46
timeformat(transaction_time, "DAY"), Monday
where transaction_time is 12/19/2016 12:22:44
[Link] Year(date)
year(date) function takes as an argument any one of the following data types: time interval or
time stamp. The returned value represents the year of the instant in time represented by this
date object. Returns a big integer indicating the year represented by this date.
Examples
Function Result
year(12/06/17 09:15:22 AM) 17
year(2015:07:21 12:45:35 PM) 2015
Note:
Only SRID 8307 is
supported in the
current release.
5-21
Chapter 5
Applying Functions to Create a New Column
Function Result
CreatePoint(78995333342435,-122.4005650 point
002481937,8307)
Function Result
distance(37.78371333337545, 1.1394718018250743E7
-122.4052500001069, 37.78371333337545,
37.78371333337545, 8307)
5-22
Chapter 5
Applying Functions to Create a New Column
Function Result
dsintervaltonum(calc_7, "MINUTE") 301.0
,
dsintervaltonum(calc_7, "DAY") 5.016666666666667
Function Result
numtodsinterval(34, "MONTH") 2 yy 10 mm
numtodsinterval(26.5, "HOUR") 1 dd 2 hr 30 mm 0 sec
5-23
Chapter 5
Applying Functions to Create a New Column
Function Result
numtodsinterval(1230, "MINUTE") 00 dd 20 hr 30 mm 0 sec
numtodsinterval(1, "DAY") 1 dd 0 hr 0 mm 0 sec
Function Result
numtoyminterval(10.5, "YEAR") 10 yy 6 mm
numtoyminterval(34, "MONTH") 2 yy 10 mm
[Link] to_dsinterval(value1)
Function converts a string in format 'DD HH:MM:SS' into a INTERVAL DAY TO SECOND data
type. The DD part indicates the number of days between 0 to 99. The HH:MM:SS part
indicates the number of hours, minutes and seconds in the interval from 0:0:0 to
23:59:59.999999. The second part can accept upto 6 decimal places.
Input value can be one of the following types: text.
Returned value will be of type interval.
Function Result
to_dsinterval("02 23:34:12") 2 dd 23hr 34mm 12 sec
[Link] to_yminterval(value1)
Function converts a string in format 'YY-MM' into a INTERVAL YEAR TO MONTH data
[Link] YY part indicates the number of years between 0 to 99. The MM part indicates the
number of months between 0-11.
Value can be one of the following types: text.
Returned value type will be intervalym.
Function Result
to_yminterval("94-3") 94 yy 3 mm
5-24
Chapter 5
Applying Functions to Create a New Column
Function Result
ymintervaltonum(94yy 5mm,"MONTH') 1133.0
5-25
Chapter 5
Applying Functions to Create a New Column
5-26
Chapter 5
Applying Functions to Create a New Column
Function Result
IEEEremainder(8809,8808) -1.0
[Link] abs(value1)
Returns the Absolute value of the input argument.
Input value can be one of the following types: number, big integer, double, integer, float.
Returned value type will be the same as the input argument type.
Function Result
abs(1234.560789) 1234.56078
abs(0.67) 0.6700000166893005
[Link] acos(value1)
Returns the Arc cosine of a value.
Value can be one of the following types: big integer, double, integer, float.
Returned value will be of type double.
Function Result
acos(0.5) 1.0471975511965979
[Link] asin(value1)
Computes the arc sine of a value.
Value 1 can be one of the following types: big integer, double, integer, float.
Returned value type will be double.
5-27
Chapter 5
Applying Functions to Create a New Column
Function Result
asin(0.5) 0.5235987755982989
[Link] atan(value1)
Returns the arc tangent of the input value.
Input value can be one of the following types: big integer, double, integer, float.
Returned value will be of type double.
Function Result
atan(34) 1.5413930385908916
[Link] atan2
Returns the polar angle of a point (value2, value1).
Value 1 can be one of the following types: big integer, double, integer, float.
Value 2 can be one of the following types: big integer, double, integer, float.
Returned value will be of type double.
Function Result
atan2(8681.44, 8682.44) 0.7853405725825559
Function Result
binomial(8609.4, 38) 5.955734227594785E104
Function Result
bitMaskWithBitsSetFromTo(23, 23) 8388608.0
5-28
Chapter 5
Applying Functions to Create a New Column
[Link] cbrt()
Returns the cubic root of a value.
Value 1 can be one of the following types: big integer, double, integer, float.
Returned value type will be double.
Function Result
cbrt(27) 3
[Link] ceil()
Round to ceiling.
The input arguments can be one of the following data types: double, float.
Returned value type will be float.
Function Result
ceil(65) 65.0
[Link] copySign()
Function returns the first floating-point argument with the sign of the second floating-point
argument.
Value1 can be one of the following types: double, float.
Returned value type will be double, float.
Function Result
copySign(3.0, -4.0)) -3.0
[Link] cos(value1)
Returns the cosine of a value
Value 1 can be one of the following types: big integer, double, integer, float.
Returned value type will be double.
Function Result
cos(7740.8) 0.9964325256163951
[Link] cosh(value1)
Returns the Cosine hyperbolic of a value.
Value 1 can be one of the following types: big integer, double, integer, float.
5-29
Chapter 5
Applying Functions to Create a New Column
Function Result
cosh(0.5) 1.1276259652063807
Function Result
exp(10) 22026.465794806718
[Link] expm1(value1)
Returns the more precise equivalent of Exp(x)-1 when x is around zero.
Value can be one of the following types: big integer, double, integer, float.
Returned value will be of type double.
Function Result
expm1(0.7) 1.0137526834646737
[Link] factorial(value1)
Returns the Factorial of a natural.
Value 1 can be one of the following types: integer.
Returned value type will be double.
Function Result
factorial(6) 720.0
[Link] floor(value1)
Value can be one of the following types: big integer, double, integer, float.
Returned value will be of type float.
Function Result
floor(0.567) 0.0
[Link] GetExponent(value1)
Function returns the unbiased exponent used in the representation of a double.
Value 1 can be one of the following types: double, float.
5-30
Chapter 5
Applying Functions to Create a New Column
Function Result
getExponent(10.0) 3.0
Function Result
getSeedAtRowColumn(48, 2) 443210610
[Link] hash(value1)
Function returns an integer hashcode for the specified value.
Value can be one of the following types: big integer, double, integer, float.
Returned value will be of type integer.
Function Result
hash(8.1) 1.33589862E9
Function Result
hypot(2,4) 4.47213595499958
[Link] LeastSignificantBit(value1)
Method is used to return the least significant 64 bits of this UUID's 128 bit value.
Value 1 can be one of the following types: integer.
Returned value will be same as the input argument.
5-31
Chapter 5
Applying Functions to Create a New Column
Function Result
LeastSignificantBit(2) 1.0
Function Result
log(20,3) 0.3667257913420846
[Link] log1(value1)
Function returns the natural logarithm of a number.
Value 1 can be one of the following types: double, integer, float.
Returned value will be of type double.
Function Result
log1(20) 2.995732273553991
[Link] log10(value1)
Logarithm(10, arg)
Value 1 can be one of the following types: big integer, double, integer, float.
Returned value will be of type double.
Function Result
log10(20) 1.301029995663981
[Link] log2(value1)
Logarithm(2, arg)
Value 1 can be one of the following types: big integer, double, integer, float.
Returned value will be of type double.
Function Result
log2(20) 4.321928094887362
5-32
Chapter 5
Applying Functions to Create a New Column
[Link] logFactorial(value1)
Function returns the natural logarithm (base e) of the factorial of its integer argument as a
double.
Value 1 can be one of the following types: double, integer, float.
Returned value will be of type double.
Function Result
logFactorial(20) 42.335616460753485
[Link] long()
Converts the input argument value to long. The input argument can be one of the following
types: big integer, integer, text, float, timestamp. Returned value type will be big integer.
Function Result
long(5039505078907524) 5039505078907524
long(22) 22
[Link] longFactorial(value1)
Function returns the natural logarithm (base e) of the factorial of its integer argument as a
double.
Value 1 can be one of the following types: double, integer, float.
Returned value will be of type double.
Function Result
longFactorial(10) 15.104412573075516
Examples
Function Result
minimum(16324, 16321) 16321
minimum(3.16, 3.10) 3.10
5-33
Chapter 5
Applying Functions to Create a New Column
Note:
If the user provides two different data types as arguments, then Stream Analytics
does implicit conversion to convert one argument to the other argument’s type.
Function Result
mod(10,3) 1.0
[Link] mostSignificantBit(value1)
Function returns the most significant 64 bits of this UUID's 128 bit value .
Value 1 can be one of the following types: integer.
Returned value will be of the same type as the first argument.
Function Result
mostSignificantBit(10) 3.0
Function Result
nextAfter()
5-34
Chapter 5
Applying Functions to Create a New Column
Function Result
nextDown()
[Link] nextUp(value1)
Function returns the floating-point number adjacent to the first argument in the direction of the
second argument.
Value 1 can be one of the following types: double, float.
Returned value will be the same type as the first argument.
Function Result
nextUp()
Function Result
pow(12,2) 144
[Link] rint(value1)
Returns the double value that is closest in value to the argument and is equal to a
mathematical integer.
Value can be one of the following types: double.
Returned value will be of type double.
Function Result
rint()
[Link] round(value1)
Rounds the argument value to the nearest integer value. The input argument can be of the
following data types: big integer, double, integer, float.
Examples
Function Result
round(7.16) 7
round(38.941) 39
5-35
Chapter 5
Applying Functions to Create a New Column
Function Result
round(3.5) 4
[Link] scalb(
Function Return d × 2scaleFactor rounded as if performed by a single correctly rounded
floating-point multiply to a member of the double value set.
Value 1 can be one of the following types: double, float.
Value 2 can be one of the following types: integer.
Returned value will be the same type as the first argument.
Function Result
scalb(10.0,2) 40.0
[Link] signum(value1)
Signum of an argument as a double value.
Value 1 can be one of the following types: number, big integer, double, integer, float.
Returned value will be of type integer.
Function Result
signum(10) 1.0
[Link] sin(value1)
Returns the sine of the input value.
Value 1 can be one of the following types: big integer, double, integer, float.
Returned value will be of type double.
Function Result
sin(7740.8) 0.08419864005868474
[Link] sinh(value1)
Returns the Sine hyperbolic of a value.
Value 1 can be one of the following types: big integer, double, integer, float.
Returned value will be of type double.
Function Result
sinh(0.5) 0.5210953054937474
5-36
Chapter 5
Applying Functions to Create a New Column
[Link] sqrt(value1)
Returns the Square root of the input value.
Value 1 can be one of the following types: big integer, double, integer, float.
Returned value will be of the type double.
Function Result
sqrt(7434.73) 86.22488040003303
[Link] stirlingCorrection(value1)
Returns the correction term of the Stirling approximation of the natural logarithm (base e) of the
factorial of the integer argument as a double: STIRLINGCORRECTION
Value 1 can be one of the following types: integer.
Returned value will be of the type double.
Function Result
stirlingCorrection(70) 0.0011904680924708464
[Link] tan(value1)
Returns the Tangent of a value.
Value 1 can be one of the following types: big integer, double, integer, float.
Returned value will be of type double.
Function Result
tan(60) 0.320040389379563
[Link] tanh(value1)
Returns the Tangent hyperbolic of a value.
Value 1 can be one of the following types: big integer, double, integer, float.
Returned value will be of type double.
Function Result
tanh(1) 0.7615941559557649
[Link] toDegrees(value1)
Converts the argument value to degrees. The input argument is an angle in radians and can be
of type double. The returned value will be the measurement of the angle in degrees and is of
type double.
5-37
Chapter 5
Applying Functions to Create a New Column
Examples
Function Result
toDegrees(3.14) 180.0
toDegrees(0.785) 45.0
[Link] toRadians(value1)
Converts the argument value to radians. The input argument is an angle in degrees and can be
of type double. The returned value will be the measurement of the angle in radians and is of
type double.
Examples
Function Result
toRadians(180.0) 3.14
toRadians(45.0) 0.785
[Link] ulp(value1)
Returns the returns the size of an ulp of the argument: ULP.
The input arguments can be one of the following data types: double, float.
Returned value type will be the same as the input value type.
Function Result
ulp(1451.54) 2.2737367544323206E-13
5-38
Chapter 5
Applying Functions to Create a New Column
Example
Function Result
nvl(Not Applicable,Commission) Not Applicable
5-39
Chapter 5
Applying Functions to Create a New Column
Function Result
beta1(0.1, 1.1, 0.2) 0.8620112116492348
beta1(316.13, 316.13, 0.2) 1.40801423421089E-63
Function Result
betacomplemented(0.1, 1.1, 0.2) 0.017407170120127144
5-40
Chapter 5
Applying Functions to Create a New Column
Function Result
beta1(2, 2, 0.5) 1.0
Function Result
binomialcomplemented(2, 3, 0.5) 0.125
Function Result
chiSquare(3.0, 5.0) 0.8282028557032665
Function Result
chiSquareComplemented(value1, value2) 0.1717971442967335
[Link] errorFunction(value1)
Returns the error function of the normal distribution.
Value 1 can be one of the following types: double, float.
5-41
Chapter 5
Applying Functions to Create a New Column
Function Result
errorFunction(5.0) 0.9999999999984626
[Link] errorFunctionComplemented(value1)
Returns the complementary Error function of the normal distribution.
Value 1 can be one of the following types: double, float.
Returned value type will be double.
Function Result
errorFunctionComplemented(5.0) 1.5374597944280347E-12
Function Result
gamma(1.0,2.0,5.0) 0.04042768199451279
Function Result
gammacomplemented(1.0, 2.0, 5.0) 0.04042768199451279
5-42
Chapter 5
Applying Functions to Create a New Column
Function Result
incompleteBeta(1.0,2.0,0.5) 0.75
Function Result
incompleteGamma(1.0,2.0) 0.8646647167633873
Function Result
incompleteGammaComplement(1.0, 2.0) 0.1353352832366127
[Link] logGamma(value1)
Returns the natural logarithm of the gamma function
Value can be one of the following types: double, float.
5-43
Chapter 5
Applying Functions to Create a New Column
Function Result
logGamma(7795.6) 62059.66356433673
Function Result
negativeBinomial(1,2,0.5) 0.5
Function Result
negativeBinomialComplemented(1.0, 2.0, 0.5
0.5)
5-44
Chapter 5
Applying Functions to Create a New Column
Function Result
normal(5.0,3.0,0.5) 0.004687384229717484
[Link] normalInverse(value1)
Returns the value for which the area under the Normal (Gaussian) probability density function
is equal to the argument value1 (assumes mean is zero, variance is one).
Input value should be between 0 and 1. The value can be one of the following types: double,
float. Required.
Returned value type will be double.
Function Result
normalInverse(0.5) 0.0
Function Result
poisson(61,123.75) 3.0509714140892473E-10
Function Result
poissonComplemented(5,3.0) 0.08391794203130347
5-45
Chapter 5
Applying Functions to Create a New Column
Value2: The integration end point. Value can be one of the following types: double, float.
Required.
Returned value type will be double.
Function Result
studentT(2.0,5.0) 0.9811252243246882
Function Result
studentTInverse(0.5, 10) 0.6998121397488263
5-46
Chapter 5
Applying Functions to Create a New Column
[Link] coalesce(value1,... )
coalesce returns the first non-null expression in the list of expressions. You must specify at
least two expressions. If all expressions evaluate to null then the coalesce function will return
null.
For example:
In coalesce(expr1,expr2):
[Link] Concat(value1,...)
Concat(value1,...) - Concatenation of values converted to strings
Value1: A part of string to concatenate with others. Value can be one of the following types: big
integer, number, double, text, integer, float, timestamp. Required.
Vararg1: A part of string to concatenate with others. Value can be one of the following types:
big integer, number, double, text, integer, float, timestamp. Optional.
Returned value will be of type text.
Function Result
Concat(client_name, card_number) Declan BENNETT0142354466948788
5-47
Chapter 5
Applying Functions to Create a New Column
Function Result
indexof(client_name,"c"), where client name 17
is Alphonse Gabriel Capone
indexof(client_name,"c"), where client name -1
is Braden Gray
[Link] initcap(value1)
Function returns a specified text expression, with the first letter of each word in uppercase and
all other letters in lowercase : INITCAP
Value1: A text expression. Value can be one of the following types: text.
Returned value will be of type text.
Function Result
initcap(client_name), where client name is Owen Taylor
Owen TAYLOR
[Link] length(value1)
Returns the length in characters of the string passed as an input argument. The input
argument is of the data type text. The returned value is an integer representing the total length
of the string.
If value1 is null, then length(value1) returns null.
Examples
Function Result
length(“one”) 3
length() ERROR: Function has invalid parameters.
length(“john”) 4
length(” “) NULL
length(null) NULL
length(“[Link]@[Link]” 30
)
5-48
Chapter 5
Applying Functions to Create a New Column
Function Result
like(client_name, "ADAMS"), True
where client name is Cameron Adams
like(client_name, "ADAMS"), False
where client name is Levi Gray
[Link] lower(value1)
Converts a string to all lower-case characters. The input argument is of the data type text. The
returned value is the lowercase of the specified string.
Examples
Function Result
lower(“PRODUCT”) product
lower(“ABCdef”) abcdef
lower(“abc”) abc
Function Result
lpad("David",10,"e") eeeeeDavid
Function Result
ltrim(client_name, "A"), where client_name lphonse Gabriel CAPONE
is Alphonse Gabriel CAPONE
5-49
Chapter 5
Applying Functions to Create a New Column
If match is not found in the string, then the original string will be returned.
Examples
Function Result
replace(“aabbccdd”,”cc”,”ff”) aabbffdd
replace(“aabbcccdd”,”cc”,”ff”) aabbffcdd
replace(“aabbddee”,”cc”,”ff”) aabbddee
Function Result
rpad("Levi Cruz", 25, "a") Levi Cruzaaaaaaaaaaaaaaaa
Function Result
rtrim(client_name, "S"), where client_name Cooper DAVI
is Cooper DAVIS
5-50
Chapter 5
Applying Functions to Create a New Column
[Link] substr()
Substr(string, from) - Substring of a 'string' when indices are between 'from' (inclusive) and up
to the end of the string.
Value 1 can be one of the following types: text.
Value 2 can be one of the following types: integer.
Returned value type will be text.
Function Result
substr(client_name, 4),where client_name is n THOMPSON
Logan THOMPSON
Examples
Function Result
substring(“abcdefgh”,3,7) cdef
substring(“abcdefgh”,1,6) abcde
Function Result
translate(client_name, "JONES", Cooper Mark
"Mark"), where the value for
client_name is Cooper JONES.
5-51
Chapter 5
Adding Custom Functions and Custom Stages
[Link] upper(value1)
Converts a string to all upper-case characters. The input argument is of the data type text. The
returned value is the uppercase of the specified string.
Examples
Function Result
upper(“name”) NAME
upper(“abcdEFGH”) ABCDEFGH
upper(“ABCD”) ABCD
5-52
Chapter 5
Adding Custom Functions and Custom Stages
Note:
Functions with same name within same package/class/method in same/different jar
are not supported.
package [Link];
import [Link];
import [Link];
import [Link];
try {
MessageDigest md = [Link]("MD5");
[Link]([Link]());
byte[] digest = [Link]();
result =
[Link](digest);
} catch (NoSuchAlgorithmException e) {
[Link]();
}
return result;
}
5-53
Chapter 5
Adding Custom Functions and Custom Stages
2. Right-click the stage after which you want to add a custom stage. Click Add a Stage, and
Custom, and then select Custom Stage from Custom Jars.
3. Enter a name and suitable description for the custom stage and click Save.
4. In the stage editor, enter the following details:
a. Custom Stage Type: Select the custom stage that was previously installed though a
custom jar
b. Input Mapping: Select the corresponding column from the previous stage for every
input parameter
You can add multiple custom stages based on your use case.
package [Link];
import [Link].*;
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
@SuppressWarnings("serial")
@OsaStage(name = "md5", description = "Create an md5 hex from a string",
inputSpec = "input, message:string", outputSpec = "output, message:string,
md5:string")
public class CustomMD5Stage implements EventProcessor {
EventFactory eventFactory;
EventSpec outputSpec;
@Override
public void init(ProcessorContext ctx, Map<String, String> config) {
eventFactory = [Link]();
OsaStage meta = [Link]([Link]);
String spec = [Link]();
outputSpec = [Link](spec);
}
@Override
public void close() {
}
@Override
public Event processEvent(Event event) {
Attr attr = [Link]("message");
Map<String, Object> values = new HashMap<String, Object>();
if (![Link]()) {
String val = (String) [Link]();
5-54
Chapter 5
Adding Custom Functions and Custom Stages
import [Link].*;
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
class BookResult {
String isbn;
String title;
String publishedDate;
String publisher;
}
5-55
Chapter 5
Adding Custom Functions and Custom Stages
@SuppressWarnings("serial")
@OsaStage(name = "RestBooks", description = "Provide info for a given book",
inputSpec = "input, isbn:string", outputSpec = "output, isbn:string,
title:string, publishedDate:string, publisher:string")
public class CustomStageRest implements EventProcessor {
EventFactory eventFactory;
EventSpec outputSpec;
static {
try {
[Link]([Link]("/
[Link]"));
} catch (IOException ioex) {
[Link]();
}
}
@Override
public void init(ProcessorContext ctx, Map<String, String> config) {
eventFactory = [Link]();
OsaStage meta = [Link]([Link]);
String spec = [Link]();
outputSpec = [Link](spec);
}
@Override
public void close() {
}
@Override
public Event processEvent(Event event) {
Attr isbnAttr = [Link]("isbn");
[Link]("isbn", isbn);
[Link]("title", [Link]);
[Link]("publishedDate", [Link]);
[Link]("publisher", [Link]);
} else {
[Link]("isbn", "");
[Link]("title", "");
[Link]("publishedDate", "");
[Link]("publisher", "");
}
Event outputEvent = [Link](outputSpec, values,
[Link]());
5-56
Chapter 5
Adding Custom Functions and Custom Stages
return outputEvent;
}
/**
* Calls the Google Books REST API to get book information based on the
ISBN ID
* @param isbn
* @return BookResult book information
*/
public BookResult getBook(String isbn) {
HttpRequestBase request;
BookResult result = null;
try {
HttpResponse response = [Link](request);
String resultJson = [Link]([Link]());
StatusLine sl = [Link]();
int code = [Link]();
if (code < 200 || code >= 300) {
[Link]("" + code + " : " + [Link]());
}
if ([Link]() > 0) {
result = new BookResult();
JsonNode book = [Link](0).path("volumeInfo"); // We
only consider the first book for this ISBN
[Link] = isbn;
[Link] = [Link]("title").asText();
[Link] = [Link]("publishedDate").asText();
[Link] = [Link]("publisher").asText();
return result;
} else {
return null; // No book found
}
5-57
Chapter 5
Adding Custom Functions and Custom Stages
} catch (Exception e) {
[Link]();
return null;
}
}
}
Note:
Following third-party jars are required for compilation of REST sample,
• [Link]
• [Link]
• [Link]
The above jars are required only at compile time and need not be packaged along
with custom jar. These libraries and their dependencies are already packaged with
OSA distribution.
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
import [Link];
@SuppressWarnings("serial")
@OsaStage(name = "CustomSoapBatchCall", description = "Call a Hello World
Soap WS", inputSpec = "input, message:string", outputSpec = "output,
message:string, result:string")
public class CustomSoapBatchCall implements BatchEventProcessor {
EventFactory eventFactory;
EventSpec outputSpec;
URL url;
QName qname;
Service service;
HelloWorldServer server;
5-58
Chapter 5
Adding Custom Functions and Custom Stages
@Override
public void init(ProcessorContext ctx, Map<String, String> config) {
eventFactory = [Link]();
OsaStage meta =
[Link]([Link]);
String spec = [Link]();
outputSpec = [Link](spec);
try {
url = new URL("[Link]
} catch (MalformedURLException e) {
[Link]();
}
qname = new QName("[Link]
"HelloWorldServerImplService");
service = [Link](url, qname);
server = (HelloWorldServer) [Link]([Link]);
}
@Override
public void close() {
}
@Override
public Iterator<Event> processEvents(Iterator<Event> iterator) {
List<String> reqs = new ArrayList<String>();
while([Link]()){
[Link]((String)[Link]().getAttr("message").getObjectValue());
@Override
public Event next() {
Map<String, Object> values = new HashMap<String, Object>();
[Link]("message",[Link](i-1));
[Link]("result",ress[i-1]);
return
[Link](outputSpec,values,[Link]());
};
5-59
Chapter 5
Adding Custom Functions and Custom Stages
@Override
public Event processEvent(Event event) {
Attr attr = [Link]("message");
Map<String, Object> values = new HashMap<String, Object>();
if (![Link]()) {
String val = (String) [Link]();
String result = callSoap(val);
[Link]("message", val);
[Link]("result", result);
} else {
[Link]("message", "empty");
[Link]("result", "empty");
}
Event outputEvent = [Link](outputSpec,
values,[Link]());
return outputEvent;
}
5.5.5 Limitations
The limitations and restrictions of the custom stages and custom functions are listed in this
section.
Custom stage type and custom functions must:
• only be used for stateless transformations. Access to state from previous calls to stage
type or function methods cannot be guaranteed and might change based on optimizations.
• not use any blocking invocations.
• not start a new thread.
• not use any thread synchronization primitives, including the wait() method, which could
potentially introduce deadlocks.
• have/be in a fully-qualified class name.
When you use the custom stages or custom functions, be careful about the heap space usage.
Note:
The resulting jar must include all the required dependencies and third-party classes
and the size of the jar file must be less than 160 MB.
5-60
Chapter 5
Writing CQL Queries
Sample: Data_A_Followed_B
5-61
Chapter 5
Writing CQL Queries
11 6 BOOKED 76547
12 2 SHIPPED 4345.98
Sample: Change_Detector
[Link] A Followed By B
Note:
In the sample query below, update q1 to the correct name of the previous stage:
Replace FROM q1 to FROM <previous-stage-name>.
5-62
Chapter 5
Writing CQL Queries
SELECT
order_id AS order_id,
abInterval AS abInterval,
Trans_id AS Trans_id,
aState_Trans_id AS aState_Trans_id,
order_status AS order_status,
aState_order_status AS aState_order_status,
order_revenue AS order_revenue,
aState_order_revenue AS aState_order_revenue
FROM q1
MATCH_RECOGNIZE (
PARTITION BY
order_id
MEASURES
B.Trans_id AS Trans_id,
B.order_id AS order_id,
B.order_status AS order_status,
B.order_revenue AS order_revenue,
A.Trans_id AS aState_Trans_id,
A.order_status AS aState_order_status,
A.order_revenue AS aState_order_revenue,
(to_timestamp(B.ELEMENT_TIME) - to_timestamp(A.ELEMENT_TIME)) AS abInterval
PATTERN( A C*? B )
WITHIN 1 minutes
DEFINE
A as A.order_status like ".*BOOKED.*" ,
B as B.order_status like ".*SHIPPED.*"
) as M
Output
{"order_id":1,"abInterval":"+000000000
00:00:05.000000000","Trans_id":6,"aState_Trans_id":1,"order_status":"SHIPPED",
"aState_order_status":"BOOKED","order_revenue":4345.0,"aState_order_revenue":2
345.98}
{"order_id":3,"abInterval":"+000000000
00:00:05.000000000","Trans_id":8,"aState_Trans_id":3,"order_status":"SHIPPED",
"aState_order_status":"BOOKED","order_revenue":3468.87,"aState_order_revenue":
3468.87}
{"order_id":2,"abInterval":"+000000000
00:00:10.000000000","Trans_id":12,"aState_Trans_id":2,"order_status":"SHIPPED"
,"aState_order_status":"BOOKED","order_revenue":4345.98,"aState_order_revenue"
:4345.98}
5-63
Chapter 5
Writing CQL Queries
Note:
In the sample query below, update q1 to the correct name of the previous stage:
Replace FROM q1 to FROM <previous-stage-name>.
SELECT
Trans_id AS Trans_id,
order_id AS order_id,
order_status AS order_status,
order_revenue AS order_revenue
FROM q1
MATCH_RECOGNIZE (
PARTITION BY
order_id
MEASURES
A.Trans_id AS Trans_id,
A.order_id AS order_id,
A.order_status AS order_status,
A.order_revenue AS order_revenue
INCLUDE TIMER EVENTS
PATTERN( A B* )
DURATION 1 minutes
DEFINE
A as A.order_status like ".*BOOKED.*" ,
B as NOT (B.order_status like ".*SHIPPED.*" )
) as M
Output
{"Trans_id":9,"order_id":4,"order_status":"BOOKED","order_revenue":3456.0}
{"Trans_id":10,"order_id":5,"order_status":"BOOKED","order_revenue":6546.0}
{"Trans_id":11,"order_id":6,"order_status":"BOOKED","order_revenue":76547.0}
Note:
In the sample query below, update q1 to the correct name of the previous stage:
Replace FROM q1 to FROM <previous-stage-name>.
5-64
Chapter 5
Writing CQL Queries
RSTREAM( SELECT
count(*) AS Number_of_Duplicates,
eventSource.order_id AS order_id,
current(eventSource.Trans_id) AS Trans_id,
current(eventSource.order_status) AS order_status,
current(eventSource.order_revenue) AS order_revenue
FROM q1 [now] as eventSource,
q1 [range 1 minutes] as dup
WHERE
eventSource.order_id = dup.order_id
GROUP BY
eventSource.order_id HAVING count(*) > 1)
Output
{"Number_of_Duplicates":2,"order_id":1,"Trans_id":4,"order_status":"PAID","ord
er_revenue":2345.98}
{"Number_of_Duplicates":2,"order_id":2,"Trans_id":5,"order_status":"PAID","ord
er_revenue":4345.98}
{"Number_of_Duplicates":3,"order_id":1,"Trans_id":6,"order_status":"SHIPPED","
order_revenue":4345.0}
{"Number_of_Duplicates":2,"order_id":3,"Trans_id":7,"order_status":"PAID","ord
er_revenue":3468.87}
{"Number_of_Duplicates":3,"order_id":3,"Trans_id":8,"order_status":"SHIPPED","
order_revenue":3468.87}
{"Number_of_Duplicates":3,"order_id":2,"Trans_id":12,"order_status":"SHIPPED",
"order_revenue":4345.98}
Note:
In the sample query below, update q1 to the correct name of the previous stage:
Replace FROM q1 to FROM <previous-stage-name>.
SELECT
Stock_ID AS Stock_ID,
Stock_Price AS Stock_Price,
orig_Stock_Price AS orig_Stock_Price,
Msg_ID AS Msg_ID,
orig_Msg_ID AS orig_Msg_ID
FROM q1
MATCH_RECOGNIZE (
5-65
Chapter 5
Writing CQL Queries
PARTITION BY
Stock_ID
MEASURES
Z.Stock_ID AS Stock_ID,
last(X.Stock_Price) AS Stock_Price,
Z.Stock_Price AS orig_Stock_Price,
last(X.Msg_ID) AS Msg_ID,
Z.Msg_ID AS orig_Msg_ID
PATTERN( Z X+ )
WITHIN 1 minutes
DEFINE X as X.Stock_Price != Z.Stock_Price
) as M
Output
{"Stock_ID":1,"Stock_Price":105,"orig_Stock_Price":87,"Msg_ID":23,"orig_Msg_ID
":4}
{"Stock_ID":2,"Stock_Price":112,"orig_Stock_Price":41,"Msg_ID":27,"orig_Msg_ID
":16}
{"Stock_ID":3,"Stock_Price":176,"orig_Stock_Price":65,"Msg_ID":30,"orig_Msg_ID
":10}
Note:
In the sample query below, update q1 to the correct name of the previous stage:
Replace FROM q1 to FROM <previous-stage-name>.
ISTREAM( SELECT
Msg_ID AS Msg_ID,
Stock_ID AS Stock_ID,
Stock_Price AS Stock_Price
FROM q1 [range 1 minutes]
) DIFFERENCE USING (Stock_Price)
Output
{"Msg_ID":1,"Stock_ID":1,"Stock_Price":87}
{"Msg_ID":5,"Stock_ID":2,"Stock_Price":41}
{"Msg_ID":6,"Stock_ID":3,"Stock_Price":65}
{"Msg_ID":17,"Stock_ID":1,"Stock_Price":91}
{"Msg_ID":20,"Stock_ID":1,"Stock_Price":105}
{"Msg_ID":24,"Stock_ID":2,"Stock_Price":112}
{"Msg_ID":28,"Stock_ID":3,"Stock_Price":176}
5-66
6
Analyze
Using Geofences for Location-based Analytics
Transforming and Analyzing Data using Patterns
Using Machine Learning Models for Scoring and Prediction
Integrating with Druid Timeseries Database for Realtime Interactive Analytics
6-1
Chapter 6
Using Geofences for Location-based Analytics
6-2
Chapter 6
Using Geofences for Location-based Analytics
6-3
Chapter 6
Using Geofences for Location-based Analytics
Note:
If you do not provide the Lib Url, Oracle maps uses the default Lib Url:
[Link]
• Authentication Type:
– API key — API key for authentication which you can get from Google. The URL
would look like this: [Link]
key=YOUR_API_KEY&callback=initMap.
– Client ID — your client ID which identifies you as a Maps API for Business
customer
6. Click Save.
Note:
To use Google maps tile layer, the usage of the maps must meet the terms of service
defined by Google ([Link]
6-4
Chapter 6
Using Geofences for Location-based Analytics
6-5
Chapter 6
Using Geofences for Location-based Analytics
Note:
Once you have modified the global parameters to customize the tile layer, the map is
updated to use the custom tile layer. These customizations, will then be applied to all
geofences.
6-6
Chapter 6
Using Geofences for Location-based Analytics
Note:
You cannot edit or update database-based geo fences.
6-7
Chapter 6
Using Geofences for Location-based Analytics
6-8
Chapter 6
Using Geofences for Location-based Analytics
• Object Key: Select field that uniquely identifies object. E.g. Vehicle Id.
• Coordinate System: The default value is 8307 and this is the only value supported.
Note:
Make sure that you do not use any names for the fields that are already part of the
incoming stream.
The outgoing shape displays direction as one of the columns, which is of type String along
with the incoming shape.
6-9
Chapter 6
Using Geofences for Location-based Analytics
6-10
Chapter 6
Using Geofences for Location-based Analytics
To use this pattern, provide suitable values for the following parameters:
• Latitude: Select field containing latitude value from stream 1.
• Longitude: Select field containing longitude value from stream 1.
• Object Key: Select field that uniquely identifies object in stream 1.
• Event Stream 2: Select the second event stream.
• Latitude: Select field containing latitude value from stream 2.
• Longitude: Select field containing longitude value from stream 2.
• Object Key: Select field that uniquely identifies object in stream 2.
• Coordinate System: The default value is 8307 and this is the only value supported.
• Distance Buffer: Enter a proximity value for the distance buffer. This field acts as a filter
criteria of two objects and the objects that do not fall in this distance (distance between
them is more than chosen distance buffer) are filtered from result set. This value must be
less than 10000 kilometers.
Note:
When a pipeline with this pattern has a database reference with cache enabled, the
pattern does not display any output in the live output stream.
The outgoing shape displays distance as another column, which is the distance between two
object under consideration along with the incoming shape.
6-11
Chapter 6
Using Geofences for Location-based Analytics
For example, if you have certain stores in the city of California, you can send promotional
messages as soon as the customer comes into a proximity of 1000 meters from any of the
stores.
To use this pattern, provide suitable values for the following parameters:
• Geo Fence: Select a geo fence that you like to analyze.
• Latitude: Select the field containing latitude value.
• Longitude: Select the field containing longitude value.
• Object Key: Select the field that uniquely identifies object. Example, Vehicle Id.
• Coordinate System: The default value is 8307 and this is the only value supported.
• Distance Buffer: Enter a proximity value for the distance buffer. This field acts as a filter
criteria for events and the events that do not fall in this distance (distance between them is
more than chosen distance buffer) are filtered from result set. This value must be less than
10000 kilometers.
The outgoing shape displays distance as another column, which is the distance between the
object and geo fence under consideration along with the incoming shape.
6-12
Chapter 6
Transforming and Analyzing Data using Patterns
6-13
Chapter 6
Transforming and Analyzing Data using Patterns
Category Pattern
Spatial Proximity: Stream with Geo Fence
Geo Fence
Spatial: Speed
Interaction: Single Stream
Reverse Geo Code: Near By
Geo Code
Spatial: Point to Polygon
Interaction: Two Stream
Proximity: Two Stream
Direction
Reverse Geo Code: Near By Place
Proximity: Single Stream
Geo Filter
Filter Eliminate Duplicates
Fluctuation
State 'A' Not Followed by 'B'
Inverse W
Detect Missing Event
W
'A' Followed by 'B'
‘B’ Not Preceded by ‘A’
Delay Event
Time Window Snapshot
Row Window Snapshot
Current And Previous Pattern
Finance Inverse W
W
Shape Detector Inverse W
W
Trend 'A' Not Followed by 'B'
Top N
Change Detector
Up Trend
Detect Missing Event
Down Trend
'A' Followed by 'B'
Detect Duplicates
Bottom N
Machine Learning Oracle Machine Learning Service
Statistical Correlation
Quantile
Transform ToJson
Split
6-14
Chapter 6
Transforming and Analyzing Data using Patterns
Outgoing Shape
The outgoing shape is the same as incoming shape. If there are no missing heartbeats, no
events are output. If there is a missing heartbeat, the previous event, which was used to
calculate the heartbeat interval is output.
6-15
Chapter 6
Transforming and Analyzing Data using Patterns
To use this pattern, provide suitable values for the following parameters:
• Partition Criteria: Select a field as the partition criteria.
• Observable Parameter: Select as field as the parameter to calculate the quantile.
• Phi-quantile: Select the percentile value to calculate the quantile of the selected event
stream. Values can only be from 1 to 99.
• Window: Select the range that determines the amount of data to consider.
• Slide: Select the frequency for newly updated output to be pushed downstream and into
the browser.
The outgoing shape is the same as the incoming shape.
6-16
Chapter 6
Transforming and Analyzing Data using Patterns
Outgoing Shape
The outgoing shape is the same as the incoming shape with one extra field:
Number_of_Duplicates. This extra field will carry the number of duplicate events that have
been discovered. All the other fields will have values of the last duplicate event.
6-17
Chapter 6
Transforming and Analyzing Data using Patterns
• Change Criteria: Select a list of fields to be compared. If the fields contain no changes, no
alerts will be generated.
• Alert on group changes: Select this option default group changes support. If it is OFF,
then alert on at least one field changes. If it is ON, then sends alert on every field change.
Outgoing Shape
The outgoing shape is based on the incoming shape, the difference being that all the fields
except the one in the partition criteria parameter will be duplicated to carry both the initial event
values and the change event values.
Example:
Your incoming event contains the following fields:
• sensor_id
• temperature
• pressure
• location
Normally, you would use sensor_id to partition your data, to look for changes in temperature.
So, select sensor_id in the partition criteria parameter and temperature in the change criteria
parameter. Use a range window that fits your use case. In this scenario, you will have the
following outgoing shape:
• sensor_id
• temperature
• orig_temperature
• pressure
• orig_pressure
• location
• orig_location
The orig_ fields carry values from the initial event. In this scenario, temperature and
orig_temperature values are different, while pressure and orig_pressure, location, and
orig_location may have identical values.
6-18
Chapter 6
Transforming and Analyzing Data using Patterns
To use this pattern, provide suitable values for the following parameters:
• Partition Criteria: Select the fields to be used as partition criteria. In the order example
above, it may be order_id.
• State A: field: Select an initial state field, whose value will be used in the comparison of
two events. In our example, it will be order_status.
• State A: value: Select the initial field state value. In our example, BOOKED.
• State B: field: Select a consecutive state field, whose value will be used in the comparison
of two events. In our example, it will be order_status again.
• State B: value: Select the consecutive field state value. In our example, SHIPPED.
• Duration: Select the time period within which to look for state changes.
Outgoing Shape
The outgoing shape is based on the incoming shape. A new abInterval field is added to carry
the value of the time interval between the states in nanosecond. Also, all but the partition
criteria fields are duplicated to carry values from both a and b states. For example, if you have
the following incoming shape:
• order_id
• order_status
• order_revenue
You will get the following outgoing shape:
• order_id
• abInterval
• order_status (this is the value by which you partition your stream)
• aState_order_status (this is the value of order_status in state A, in our example
'BOOKED')
• order_revenue (this is the value of order_revenue in state B)
• aState_order_revenue (this is the value of order_revenue in state A)
6-19
Chapter 6
Transforming and Analyzing Data using Patterns
To use this pattern, provide suitable values for the following parameters:
• Window Range: Select a rolling time period within which the events will be collected and
ordered per your ordering criteria.
• Window Slide: Select the frequency for the newly updated output to be pushed
downstream and into the browser.
• Order by Criteria: Select a list of fields to use to order the collection of events.
• Number of Events: Select the number of top value events to output.
The outgoing shape is the same as the incoming shape.
6-20
Chapter 6
Transforming and Analyzing Data using Patterns
Outgoing Shape
The outgoing shape is based on the incoming shape with an addition of two new fields. For
example, if your incoming event contains the following fields:
• sensor_id
• temperature
• pressure
• location
Normally, you would use sensor_id to partition your data and say you want to look for the
upward trend in temperature. So, select sensor_id in the partition criteria parameter and
temperature in the tracking value parameter. Use a duration that fits your use case. In this
scenario, you will have the following outgoing shape:
• sensor_id
• startValue (this is the value of temperature that starts the trend)
• endValue (this is the value of temperature that ends the trend)
• temperature (the value of the last event)
• pressure (the value of the last event)
• location (the value of the last event)
Outgoing Shape
The outgoing shape is based on the incoming shape with an addition of two new fields. Let's
look at an example. Your incoming event contains the following fields:
• sensor_id
• temperature
• pressure
6-21
Chapter 6
Transforming and Analyzing Data using Patterns
• location
Normally, you would use sensor_id to partition your data and say you want to look for the
downward trend in temperature. So, select sensor_id in the partition criteria parameter and
temperature in the tracking value parameter. Use a duration that fits your use case. In this
scenario, you will have the following outgoing shape:
• sensor_id
• startValue (this is the value of temperature that starts the trend)
• endValue (this is the value of temperature that ends the trend)
• temperature (the value of the last event)
• pressure (the value of the last event)
• location (the value of the last event)
The pattern is visually represented based on the data you have entered/selected.
Outgoing Shape
The outgoing shape is the same as incoming shape. If the second (state B) event does not
arrive within the specified time window, the first (state A) event is pushed to the output.
6-22
Chapter 6
Transforming and Analyzing Data using Patterns
To use this pattern, provide suitable values for the following parameters:
• Partition Criteria: (Optional) a field to partition your stream by. In the order example
above, it may be order_id.
• State A: field: an initial state field, whose value will be used in the comparison of two
events. In our example, it will be order_status.
• State A: value: the initial field state value. In our example, BOOKED.
• State B: field: a consecutive state field, whose value will be used in the comparison of two
events. In our example, it will be order_status again.
• State B: value: the consecutive field state value. In our example, SHIPPED.
• Duration: the time period, within which to look for state changes.
Outgoing Shape
The outgoing shape is the same as incoming shape. If the second (state B) event does not
arrive within the specified time window, the first (state A) event is pushed to the output.
Outgoing Shape
The outgoing shape is based on the incoming shape with an addition of five new fields. The
new fields are:
• firstW
• firstValleyW
• headW
• secondValleyW
• lastW
The new fields correspond to the tracking value terminal points of the W shape discovered in
the feed. The original fields correspond to the last event in the W pattern.
6-23
Chapter 6
Transforming and Analyzing Data using Patterns
To use this pattern, provide suitable values for the following parameters:
• Partition Criteria: Select a field to be used as a partition criterion. For example, a ticker
symbol.
• Window: Select a time period within which the values of the designated field are analyzed
for the inverse W shape.
• Tracking value: Select a field to be analyzed for the inverse W shape.
Outgoing Shape
The outgoing shape is based on the incoming shape with an addition of five new fields. The
new fields are:
• firstW
• firstPeakW
• headInverseW
• secondpeakW
• lastW
The new fields correspond to the tracking value terminal points of the inverse W shape
discovered in the feed. The original fields correspond to the last event in the inverse W pattern.
Output schema/payload from this pattern for non-partitioned input is [DiskID, Usage,
PREV_DiskId, PREV_Usage].
For input partitioned by DiskID, the output schema from the pattern is [DiskID, Usage,
PREV_Usage].
Output [Disk1, 45gb, 40gb]. There is no output for Disk2 until another event for Disk2
arrives.
6-24
Chapter 6
Transforming and Analyzing Data using Patterns
Note:
Window will automatically dump contents either when a new event arrives or when an
event expires.
6-25
Chapter 6
Transforming and Analyzing Data using Patterns
6-26
Chapter 6
Transforming and Analyzing Data using Patterns
If the size is configured to the value greater than 0 (say n), it will transform maximum n
events as a JSON array of single JSON document and output as a JSON text.
If the size is set to the default 0, all the events of a partition within the batch duration will be
transformed as an array of single JSON document and output as JSON text.
• Upload Json File: You can upload a sample JSON file to be used to infer JSON path for
field mapping.
• Field Mapping:
– Json Path: Lists all the paths in uploaded JSON file.
– Fields: Lists all the fields from the previous stage. You can map the JSON path with
one of the fields from the drop-down list.
6.2.27 Applying OML Models to get the Scoring of Events (Preview Feature)
Use the Oracle Machine Learning Service pattern to use OML models to apply the scoring
on the ingested events.
To use this pattern, provide suitable values for the following parameters:
• OML server url: Enter the OML service endpoint where the autonomous data warehouse
is located for the Machine Learning model in the region.
• Tenant: Enter the tenant ID hosting the OML model.
6-27
Chapter 6
Transforming and Analyzing Data using Patterns
6-28
Chapter 6
Using Machine Learning Models for Scoring and Prediction
• On the Predictive Model Details, enter the following details and click Save:
1. For Predictive Model URL, upload your ONNX file.
2. In the Model Version field, enter the version of this artifact. For example, 1.0.
3. (Optional) In the Version Description, enter a meaningful description for your model.
4. In the Algorithm field, accept the default. The algorithm is derived from the model you
have uploaded.
5. (Optional) In the Tool drop-down list, select the tool with which you created your
model.
6-29
Chapter 6
Integrating with Druid Timeseries Database for Realtime Interactive Analytics
6-30
Chapter 6
Integrating with Druid Timeseries Database for Realtime Interactive Analytics
• Name
• Description
• Tags
• Source Type: Select Published Pipeline from the drop-down list.
3. Click Next.
4. On the Ingestion Details screen, enter the following details:
• Connection: Select a connection from the drop-down list.
• Pipeline: Select a pipeline from the drop-down list.
• Kafka Target Select a Kafka target from the drop-down list.
• Timestamp: Select a column from the pipeline to be used as the timestamp.
• Timestamp format: Select or set a suitable format for the timestamp using Joda time
format. This is a mandatory field. The default value is auto.
• Metrics: Select metrics for creating measures.
• Dimensions: Select dimensions for group by.
• High Cardinality Dimensions: Select high cardinality dimensions such as unique IDs.
Hyperlog approximation will be used.
5. Click Next.
6. Select the required values for the metric on the Metric Capabilities screen.
7. On the Advanced Settings screen, enter the following details:
• Segment granularity: Select the granularity with which you want to create segments
• Query granularity: Select the minimum granularity to be able to query results and the
granularity of the data inside the segment
• Task count: Select the maximum number of reading tasks in a replica set. This means
that the maximum number of reading tasks is taskCount*replicas and the total
number of tasks (reading + publishing) is higher than this. The number of reading tasks
is less than taskCount if taskCount > {numKafkaPartitions}.
• Task duration: Select the length of time before tasks stop reading and begin
publishing their segment. The segments are only pushed to deep storage and loadable
by historical nodes when the indexing task completes.
• Maximum rows in memory: Enter a number greater than or equal to 0. This number
indicates the number of rows to aggregate before persisting. This number is the post-
aggregation rows, so it is not equivalent to the number of input events, but the number
of aggregated rows that those events result in. This is used to manage the required
JVM heap size. Maximum heap memory usage for indexing scales with
maxRowsInMemory*(2 + maxPendingPersists).
• Maximum rows per segment: Enter a number greater than or equal to 0. This is the
number of rows to aggregate into a segment; this number is post-aggregation rows.
• Immediate Persist Period: Select the period that determines the rate at which
intermediate persists occur. This allows the data cube is ready for query earlier before
the indexing task finishes.
• Report Parse Exception: Select this option to throw exceptions encountered during
parsing and halt ingestion.
6-31
Chapter 6
Integrating with Druid Timeseries Database for Realtime Interactive Analytics
6-32
7
Visualize
Visualizations are graphical representations of the streaming data in a pipeline. You can add
visualizations on all the stages in a pipeline, except on a target stage.
Adding Realtime Charts
Creating and Managing Dashboards
7-1
Chapter 7
Adding Realtime Charts
7-2
Chapter 7
Adding Realtime Charts
• Y Axis Field Selection: the column to be used as the Y axis. This is a mandatory
field.
• Axis Label: a label for the Y axis. This is an optional field.
• X Axis Field Selection: the column to be used as the X axis. This is a mandatory
field.
• Axis Label: a label for the X axis. This is an optional field.
• Bubble Size Field Selection: select the field that you want to use as the bubble size.
This is a mandatory field.
5. Click Create.
The visualization is created and you can see the data populated in it.
7-3
Chapter 7
Adding Realtime Charts
7-4
Chapter 7
Adding Realtime Charts
5. Click Create.
The visualization is created and you can see the data populated in it.
7-5
Chapter 7
Creating and Managing Dashboards
• Map Type: the map of the region that you want to use. This is a mandatory field.
• Location Field: the field that you want to use as the location. This is a mandatory
field.
• Data Field: the field that you want to use as the data field. This is a mandatory field.
• Show Data Value: select this check box if you want to display the data value as
marker on the visualization. This is an optional field.
5. Click Create.
The visualization is created and you can see the data populated in it.
Edit Visualization
To edit a visualization:
1. On the stage that has visualizations, click the Visualizations tab.
2. Identify the visualization that you want to edit and click the pencil icon next to the
visualization name.
3. In the Edit Visualization dialog box that appears, make the changes you want. You can
even change the Y Axis and X Axis selections. When you change the Y Axis and X Axis
values, you will notice a difference in the visualization as the basis on which the graph is
plotted has changed.
The following are the other updates you can make to the visualizations:
• Maximize Visualizations: You can open the visualization in a new window/tab using the
Maximize Visualizations icon in the visualization canvas.
• Change Orientation: Based on the data that you have in the visualization or your
requirement, you can change the orientation of the visualization. You can toggle between
horizontal and vertical orientations by clicking the Flip Chart Layout icon in the
visualization canvas.
• Delete Visualization: You can delete the visualization if you no longer need it in the
pipeline. In the visualization canvas, click the Delete icon available beside the visualization
name to delete the visualization from the pipeline. Be careful while you delete the
visualization, as it is deleted with immediate effect and there is no way to restore it once
deleted.
• Delete All Visualizations: You can delete all the visualizations in the stage if you no
longer need them. In the visualization canvas, click the Delete All icon to delete all the
visualizations of the stage at one go. Be careful while you delete the visualizations, as the
effect is immediate and there is no way to restore the deleted visualizations.
7-6
Chapter 7
Creating and Managing Dashboards
• Name
• Description
• Tags
3. Click Next.
4. On the Source Details page, enter the following details:
• Begin: Select this option to build a dashboard from Beginning or Latest offset.
Available options are latest and earliest. If you select the latest option, the dashboard
shows only the new records that have been received after opening the dashboard. If
you select the earliest option, then the records are read from the beginning of the
OSA pipeline topics.
• CSS: Enter a custom stylesheet for the dashboard.
5. Click Save.
2. Click Edit Dashboard to view the dashboard editing options under Actions.
7-7
Chapter 7
Creating and Managing Dashboards
4. Click Set autorefresh to select the refresh frequency for the dashboard. This is applicable
only for cube based visualizations not applicable for streaming charts created out of
pipeline.
This is just a client side setting and is not persisted.
Note:
The Refresh options work only for druid-based visualizations, and not for
streaming visualizations because they have continuous data flow.
5. Click the Save icon to save the changes you have made to the dashboard.
6. Click the Edit CSS icon to edit and apply a CSS to the dashboard. You can also edit the
CSS in the live editor.
You can email the dashboard link to someone using the Email the link icon.
7. Click Add visualizations to see a list of existing visualizations. Visualizations from the
pipelines and as well as from the cube explorations appear here. Go through the list, select
one or more visualizations and add them to the dashboard.
8. Hover over the added visualization, click the Explore chart icon to open the chart editor of
the visualization.
You can see the metadata of the visualization. You can also move the chart around the
canvas, refresh it, or remove it from the dashboard.
A cube exploration looks like the following:
7-8
Chapter 7
Creating and Managing Dashboards
The various options like time granularity, group by, table timestamp format, row limit, filters,
and result filters add more granularity and details to the dashboard.
7-9
Chapter 7
Creating and Managing Dashboards
7-10
8
Monitor
Execution and HA Statistics
Detailed Query Analysis
Complete CQL Engine Statistics
You can view the Execution and HA statistics, in the CQL Engine Query details page that is
displayed.
8-1
Chapter 8
Execution and HA Statistics
• Number of Partitions: This fields shows the degree of parallelism of the query. The
degree of parallelism is defined by total number of input partitions processed by a
query.
Degree of parallelism depends on the many factors such as query constructs, number
of input kafka partitions and number of executors assigned to application.
• Execution Statistics Table: This section shows the detailed execution statistics of
each operator.
– Partition ID: Partition Sequence Id
– CQL Engine ID: Sequence ID of CQL Engine on which the partition is being
processed.
– Total Output Events: Number of output events emitted by CQL query for each
partition.
– Total Output Heartbeats: Number of heartbeat events emitted by CQL query for
each partition. Please note that heartbeats are special events which ensures
timestamp progression in Oracle Stream Analytics pipeline.
– Throughput: Ratio of total number of events processed and total time spent in
processing for each partition.
– Latency: Average turnaround time taken to process a partition of stream.
• HA Statistics: This table shows the real-time statistics about query's HA operations.
Note that unit of time is in MILLISECONDS.
– Partition ID: Partition Sequence ID
– CQL Engine ID: Sequence ID of CQL Engine on which the partition is being
processed.
– Total Full Snapshots Created: Total number of times the full state of query is
serialized and saved.
– Avg Full Snapshot Creation Time: Average time spent in serializing and saving the
full state of query.
– Total Full Snapshots Loaded: Total number of times the full state of query is de-
serialized and loaded in query plan.
– Avg Full Snapshot Load Time: Average time spent in de-serializing and loading the
full state of query.
– Total Journal Snapshots Created: Total number of times the journaled state of
query is serialized and saved.
– Avg Journal Snapshot Creation Time: Average time spent in serializing and saving
the journaled state of query.
– Total Journal Snapshots Loaded: Total number of times the journaled state of
query is de-serialized and loaded in query plan.
– Avg Journal Snapshot Load Time: Average time spent in de-serializing and loading
the journaled state of query.
Full Snapshot is the complete state of query. The query state represent the internal
data structure and state of each operator in query plan. Journal snapshot is partial
and incremental snapshot having a start time and end time. Oracle Stream
Analytics optimizes the state preservation by using Journal snapshot if possible.
8-2
Chapter 8
Detailed Query Analysis
This page contains details about each execution operator of CQL query for a particular
partition of a stage in pipeline.
• Query ID: System generated identifier for query
• Query Text: Query String
• Partition ID: All operator details are corresponding to this partition id.
• Operator Statistics:
– Operator ID: System Generated Identifiers
– Total Input Events: Total number of input events received by each operator.
– Total Output Events: Total number of output events generated by each operator.
– Total Input Heartbeats: Total number of heartbeat events received by each operator.
– Total Output Heartbeats: Total number of heartbeat events generated by each
operator.
– Throughput(events/second): Ratio of total input events processed and total time spent
in processing for each operator.
– Latency(ms): Total turnaround time to process an event for each operator.
• Operator DAG: This is visual representation of the query plan. The DAG will show the
parent-child details for each operator. You can further drill down the execution statistics of
operator. Please click on the operator which will open CQL Operator Details Page.
8-3
Chapter 8
Complete CQL Engine Statistics
This page contains additional information about each execution operator, apart from the CQL
Engine Query Details page, provides all essential metrics for each operator.
Pipeline ID: Unique pipeline id in Spark Cluster
Pipeline Name: Name of Oracle Stream Analytics Pipeline given by user in Oracle Stream
Analytics UI.
Stage ID: Unique stage ID in DAG of stages for Oracle Stream Analytics Pipeline.
Running Queries: This section displays list of CQL queries running to compute the CQL
transformation for a stage. This table displays a system-generated Query ID and Query Text.
Check Oracle Continuous Query Language Reference for CQL Query syntax and semantics.
To see more details about query, click on the query id hyperlink in the table entry to open CQL
Engine Query Details page.
Registered Sources: This section displays internal CQL metadata about all the input sources
which the query is based upon. For every input stream of the stage, there will be one entry in
this table.
Each entry contains source name, source type, timestamp type and stream attributes.
Timestamp type can be PROCESSING or EVENT timestamped. If stream is PROCESSING
timestamped, then timestamp of each event will be defined by system. If stream is EVENT
timestamped, then timestamp of each event is defined by one of the stream attribute itself. A
source can be Stream or Relation.
External Sources: This section displays details about all external sources with which input
stream is joined. The external source can be a database table or coherence cache.
CQL Engines: This section displays a table having details about all instances of CQL engines
used by the pipeline. Here are details about each field of the table:
• CQLEngine Id: System generated id for a CQL engine instance.
• ExecutorId: Executor Id with which the CQL engine is associated.
• Executor Host: Address of the cluster node on which this CQL engine is running.
8-4
Chapter 8
Complete CQL Engine Statistics
• Status: Current Status of CQL Engine. Status can be either ACTIVE or INACTIVE. If it is
ACTIVE, it means that CQL Engine instance is up and running, Otherwise CQL Engine is
stopped explicitly.
8-5
9
Reference
Pipeline Details
Stage Details
Query Details
Internal Kafka Topics
9-1
Chapter 9
Stage Details
Oracle Stream Analytics supports various types of stages e.g. Query Stage, Pattern Stage,
Custom Stage etc. For each pipeline stage, Oracle Stream Analytics defines a list of
transformations which will be applied on the input stream for the stage. The output from
final transformation will be the output of stage.
• Total Output Partitions: This measurement is total number of partitions in the output
stream of each stage.
Every pipeline stage has its own partitioning requirements which are determined from
stage configurations. For example, If a QUERY stage defines a summary function and
doesn't define a group-by column, then QUERY stage will have only one partition because
there will be no partitioning criteria available.
• Total Output Events: This measurement is total number of output events (not micro-
batches) emitted by each stage.
• Average Output Rate: This measurement is the rate at which each stage has emitted
output events so far. The rate is a ratio of total number of output events so far and total
application execution time.
If the rate is ZERO, then it doesn't always mean that there is ERROR in stage processing.
Sometime stage doesn't output any record at all (e.g. No event passed the Filter in Query
stage).This can happen if rate of output events is less than 1 events/second.
• Current Output Rate: This measurement is the rate at which each stage is emitting output
events. The rate is ratio of total number of output events and total application execution
time since last metrics page refresh. To get better picture of current output rate, please
refresh the page more frequently.
page.
This page provides details about all the transformations in specific stage.
• Pipeline ID: Unique pipeline id in Spark Cluster
• Pipeline Name: Name of Oracle Stream Analytics Pipeline given by user in Oracle Stream
Analytics UI.
• Stage ID: Unique stage id in DAG of stages for Oracle Stream Analytics Pipeline.
9-2
Chapter 9
Query Details
• Stage DAG: This is a visual representation of all transformations in form of a DAG where it
displays the parent-child relation various pipeline transformations.
Note Oracle Stream Analytics allows different transformation types including Business
Rules but drilldowns are allowed only on CQLDStream transformations. Inside
CQLDStream transformation, the input data is transformed using a continuously running
query (CQL) in a CQL Engine. Note that there will be one CQL engine associated with one
Executor.
• Stage Transformations Table:This table displays details about all transformations being
performed in a specific stage. Each entry in the table is corresponding to a transformation
operation used in computation of the stage. You can observe that final transformation in
every stage is MonitorDStream. The reason is that MonitorDStream pipes output of stage
to Oracle Stream Analytics UI Live Table.
– Transformation Name: Name of output DStream for the transformation. Every
transformation in spark results into an output dstream.
– Transformation Type: This is category information of each transformation being used
in the stage execution. If transformation type is "Oracle", it is based on Oracle's
proprietary transformation algorithm. If the transformation type is "Native", then
transformation is provided by Apache Spark implementation.
CQLDStream transformation allows further drill down as described in the Query Details
section.
The CQL Engine Summary page for the query has details like query text, stream sources
feeding the query, external sources if any and all CQL engines with which the query is
registered.
9-3
Chapter 9
Internal Kafka Topics
Kafka Topics
Group IDs
9-4
10
Troubleshoot
Pipeline Debug and Monitoring Metrics
Common Issues and Remedies
You can check the status of a published application on the catalog page, as shown in the
screenshots below:
10-1
Chapter 10
Pipeline Debug and Monitoring Metrics
1. For draft pipelines, you can navigate to YARN Applications page using the YARN
Resource Manager URL and YARN Master console port values in System Settings.
2. Application master url is [Link] Resource Manager URL>:<Yarn master Console
Port>.
This page displays all the applications running in YARN:
3. After identifying the application from the list, click the ApplicationMaster link, to open the
Spark Application Details page:
10-2
Chapter 10
Pipeline Debug and Monitoring Metrics
10-3
Chapter 10
Pipeline Debug and Monitoring Metrics
• Average Output Rate: This measurement is the rate at which each stage has emitted
output events so far. The rate is a ratio of total number of output events so far and total
application execution time.
If the rate is ZERO, then it doesn't always mean that there is ERROR in stage processing.
Sometime stage doesn't output any record at all (e.g. No event passed the Filter in Query
stage).This can happen if rate of output events is less than 1 events/second.
• Current Output Rate: This measurement is the rate at which each stage is emitting output
events. The rate is ratio of total number of output events and total application execution
time since last metrics page refresh. To get better picture of current output rate, please
refresh the page more frequently.
page.
This page provides details about all the transformations in specific stage.
• Pipeline ID: Unique pipeline id in Spark Cluster
• Pipeline Name: Name of Oracle Stream Analytics Pipeline given by user in Oracle Stream
Analytics UI.
• Stage ID: Unique stage id in DAG of stages for Oracle Stream Analytics Pipeline.
• Stage DAG: This is a visual representation of all transformations in form of a DAG where it
displays the parent-child relation various pipeline transformations.
Note Oracle Stream Analytics allows different transformation types including Business
Rules but drilldowns are allowed only on CQLDStream transformations. Inside
CQLDStream transformation, the input data is transformed using a continuously running
query (CQL) in a CQL Engine. Note that there will be one CQL engine associated with one
Executor.
• Stage Transformations Table:This table displays details about all transformations being
performed in a specific stage. Each entry in the table is corresponding to a transformation
operation used in computation of the stage. You can observe that final transformation in
every stage is MonitorDStream. The reason is that MonitorDStream pipes output of stage
to Oracle Stream Analytics UI Live Table.
– Transformation Name: Name of output DStream for the transformation. Every
transformation in spark results into an output dstream.
10-4
Chapter 10
Pipeline Debug and Monitoring Metrics
The CQL Engine Summary page for the query has details like query text, stream sources
feeding the query, external sources if any and all CQL engines with which the query is
registered.
10-5
Chapter 10
Pipeline Debug and Monitoring Metrics
You can view the Execution and HA statistics, in the CQL Engine Query details page that is
displayed.
10-6
Chapter 10
Pipeline Debug and Monitoring Metrics
– Total Output Events: Number of output events emitted by CQL query for each
partition.
– Total Output Heartbeats: Number of heartbeat events emitted by CQL query for
each partition. Please note that heartbeats are special events which ensures
timestamp progression in Oracle Stream Analytics pipeline.
– Throughput: Ratio of total number of events processed and total time spent in
processing for each partition.
– Latency: Average turnaround time taken to process a partition of stream.
• HA Statistics: This table shows the real-time statistics about query's HA operations.
Note that unit of time is in MILLISECONDS.
– Partition ID: Partition Sequence ID
– CQL Engine ID: Sequence ID of CQL Engine on which the partition is being
processed.
– Total Full Snapshots Created: Total number of times the full state of query is
serialized and saved.
– Avg Full Snapshot Creation Time: Average time spent in serializing and saving the
full state of query.
– Total Full Snapshots Loaded: Total number of times the full state of query is de-
serialized and loaded in query plan.
– Avg Full Snapshot Load Time: Average time spent in de-serializing and loading the
full state of query.
– Total Journal Snapshots Created: Total number of times the journaled state of
query is serialized and saved.
– Avg Journal Snapshot Creation Time: Average time spent in serializing and saving
the journaled state of query.
– Total Journal Snapshots Loaded: Total number of times the journaled state of
query is de-serialized and loaded in query plan.
– Avg Journal Snapshot Load Time: Average time spent in de-serializing and loading
the journaled state of query.
Full Snapshot is the complete state of query. The query state represent the internal
data structure and state of each operator in query plan. Journal snapshot is partial
and incremental snapshot having a start time and end time. Oracle Stream
Analytics optimizes the state preservation by using Journal snapshot if possible.
10-7
Chapter 10
Pipeline Debug and Monitoring Metrics
This page contains details about each execution operator of CQL query for a particular
partition of a stage in pipeline.
• Query ID: System generated identifier for query
• Query Text: Query String
• Partition ID: All operator details are corresponding to this partition id.
• Operator Statistics:
– Operator ID: System Generated Identifiers
– Total Input Events: Total number of input events received by each operator.
– Total Output Events: Total number of output events generated by each operator.
– Total Input Heartbeats: Total number of heartbeat events received by each operator.
– Total Output Heartbeats: Total number of heartbeat events generated by each
operator.
– Throughput(events/second): Ratio of total input events processed and total time spent
in processing for each operator.
– Latency(ms): Total turnaround time to process an event for each operator.
• Operator DAG: This is visual representation of the query plan. The DAG will show the
parent-child details for each operator. You can further drill down the execution statistics of
operator. Please click on the operator which will open CQL Operator Details Page.
10-8
Chapter 10
Pipeline Debug and Monitoring Metrics
This page contains additional information about each execution operator, apart from the CQL
Engine Query Details page, provides all essential metrics for each operator.
Pipeline ID: Unique pipeline id in Spark Cluster
Pipeline Name: Name of Oracle Stream Analytics Pipeline given by user in Oracle Stream
Analytics UI.
Stage ID: Unique stage ID in DAG of stages for Oracle Stream Analytics Pipeline.
Running Queries: This section displays list of CQL queries running to compute the CQL
transformation for a stage. This table displays a system-generated Query ID and Query Text.
Check Oracle Continuous Query Language Reference for CQL Query syntax and semantics.
To see more details about query, click on the query id hyperlink in the table entry to open CQL
Engine Query Details page.
Registered Sources: This section displays internal CQL metadata about all the input sources
which the query is based upon. For every input stream of the stage, there will be one entry in
this table.
Each entry contains source name, source type, timestamp type and stream attributes.
Timestamp type can be PROCESSING or EVENT timestamped. If stream is PROCESSING
timestamped, then timestamp of each event will be defined by system. If stream is EVENT
timestamped, then timestamp of each event is defined by one of the stream attribute itself. A
source can be Stream or Relation.
External Sources: This section displays details about all external sources with which input
stream is joined. The external source can be a database table or coherence cache.
CQL Engines: This section displays a table having details about all instances of CQL engines
used by the pipeline. Here are details about each field of the table:
• CQLEngine Id: System generated id for a CQL engine instance.
• ExecutorId: Executor Id with which the CQL engine is associated.
• Executor Host: Address of the cluster node on which this CQL engine is running.
10-9
Chapter 10
Common Issues and Remedies
• Status: Current Status of CQL Engine. Status can be either ACTIVE or INACTIVE. If it is
ACTIVE, it means that CQL Engine instance is up and running, Otherwise CQL Engine is
stopped explicitly.
Kafka Topics
Group IDs
10-10
Chapter 10
Common Issues and Remedies
10.2.1 Pipeline
Common issues encountered while deploying pipelines are listed in this section.
10.2.2 Pipeline
Common issues encountered while deploying pipelines are listed in this section.
Ensure that the Input Stream is Supplying Continuous Stream of Events to the Pipeline
To check for a continuous supply of events from the input stream:
1. Go to the Catalog.
2. Locate and click the stream you want to troubleshoot.
3. Check the value of the topicName property under the Source Type Parameters section.
4. Since this topic is created using Kafka APIs, you cannot consume this topic with REST
APIs.
Listen to the Kafka topic hosted on a standard Apache Kafka installation.
You can listen to the Kafka topic using utilities from a Kafka Installation. kafka-console-
[Link] is a utility script available as part of any Kafka installation.
Follow these steps to listen to Kafka topic:
a. Determine the Zookeeper Address from Apache Kafka Installation based Cluster.
b. Use the following command to listen the Kafka topic:
10-11
Chapter 10
Common Issues and Remedies
The topic name is AppName_StageId. The pipeline name can be derived from topic name
by removing the _StageID from topic name. In the above snapshot, the pipeline name is
sx_2_49_12_pipe1_draft.
In case of a kafka source, republish the pipeline to read records from where it left off before
terminating.
[Link] Live Table Shows Listening Events with No Events in the Table
There can be multiple reasons why status of pipeline has not changed to Listening Events
from Starting Pipeline. Following are the steps to troubleshoot this scenario:
10-12
Chapter 10
Common Issues and Remedies
1. The live table shows output events of only the currently selected stage. At any time, only
one stage is selected. Try switching to a different stage. If you observe output in live table
for another stage, the problem can be associated with the stage. To debug further, go to
step 5.
If there is no output in any stage, then move to step 2.
2. Ensure that the pipeline is still running on Spark Cluster. See Ensure that the Pipeline is
Deployed Successfully
3. If the Spark application for your pipeline is killed or aborted, then it suggests that the
pipeline has crashed. To troubleshoot further, you may need to look into application logs.
4. If application is in ACCEPTED, NEW or SUBMITTED state, then application is waiting for
cluster resource and not yet started. If there are not enough resources, check the number
of VCORES in Big Data Cloud Service Spark Yarn cluster. For a pipeline, Stream Analytics
requires minimum 3 VCORES.
5. If application is in RUNNING state, use the following steps to troubleshoot further:
a. Ensure that the input stream is pushing events continuously to the pipeline.
b. If the input stream is pushing events, ensure that each of the pipeline stages is
processing events and providing outputs.
c. If both of the above steps are verified successfully, then ensure that the pipeline is able
to push the output events of each stage to its corresponding monitor topic:
i. Determine the monitor topic for the stage, where the output of stage is being
pushed into. See Determine the Topic Name where Output of Pipeline Stage is
Propagated .
ii. Listen to the monitor topic and ensure that the events are continuously being
pushed in topic. To listen to the Kafka topic, you must have access to Kafka cluster
where topic is created. You can listen to the Kafka topic using utilities from a Kafka
Installation. [Link] is a utility script available as part of any
Kafka installation.
iii. If you don't see any events in the topic, then this can be an issue related to writing
output events from stage to monitor topic. Check the server logs and look for any
exception and then report to the administrator.
iv. If you can see outputs events in monitor topic, then the issue can be related to
reading output events in web browser.
There can be multiple reasons why status of pipeline has not changed to Listening Events
from Starting Pipeline. Following are the steps to troubleshoot this scenario:
10-13
Chapter 10
Common Issues and Remedies
1. Ensure that the pipeline has been successfully deployed to Spark Cluster. For more
information, see Ensure that Pipeline is Deployed Successfully . Also ensure that the
Spark cluster is not down and is available.
2. If the deployment failed, check the Jetty logs to see the exceptions related to the
deployment failure and fix the issues.
3. If the deployment is successful, verify that OSA webtier has received the pipeline
deployment from Spark.
4. Click Done in the pipeline editor and then go back to Catalog and open the pipeline again.
[Link] Time-out Exception in the Spark Logs when you Unpublish a Pipeline
In the Jetty log look for the following message:
OsaSparkMessageQueue:182 - received:
[Link] Undeployment
Ack: [Link]
During an application shutdown, a pipeline may take several minutes to unpublish completely.
So, if you do not see the above message, then you may need to increase the
[Link] value accordingly.
Also in the High Availability mode, at the time of unpublishing a pipeline, the snapshot folder is
deleted.
In HA mode, if you do not receive the above error message on time, and see the following
error:
Undeployment couldn't be complete within 60000 the snapshot folder may not be
completely cleaned.
it only means that the processing will not be impacted, but some disk space will be occupied.
10-14
Chapter 10
Common Issues and Remedies
To solve this issue, check if there is a summary stage added before a target stage, and add a
query stage with a filter checking for null values for the cache keys.
10.2.3 Stream
Common issues encountered with streams are listed in this section.
[Link] Cannot See Any Kafka Topic or a Specific Topic in the List of Topics
Use the following steps to troubleshoot:
1. Go to Catalog and select the Kafka connection which you are using to create the stream.
2. Click Next to go to Connection Details tab.
3. Click Test Connection to verify that the connection is still active.
4. Ensure that topic or topics exist in Kafka cluster. You must have access to Kafka cluster
where topic is created. You can list all the Kafka topics using utilities from a Kafka
Installation. [Link] is a utility script available as part of any Kafka
installation.
5. If you can't see any topic using above command, ensure that you create the topic.
6. If the test connection failed, and you see error message like OSA-01266 Failed to
connect to the ZooKeeper server, then the Kafka cluster is not reachable. Ensure that
Kafka cluster is up and running.
[Link] Input Kafka Topic is Sending Data but No Events Seen in Live Table
This can happen if the incoming events are not adhered to the expected shape for the Stream.
To check if the events are dropped due to shape mismatch, use the following steps to
troubleshoot:
1. Verify if lenient parameter under Source Type Properties for the Stream is selected. If it is
FALSE, then the event may have been dropped due to shape mismatch. To confirm this,
check the application logs for the running application.
2. If the property is set to TRUE, debug further:
a. Make sure that Kafka Cluster is up and running. Spark cluster should be able to
access Kafka cluster.
b. If Kafka cluster is up and running, obtain the application logs for further debugging.
10.2.4 Connection
Common issues encountered with connections are listed in this section.
10-15
Chapter 10
Common Issues and Remedies
10.2.5 Target
Common issues encountered with targets are listed in this section.
10-16
Chapter 10
Common Issues and Remedies
10.2.6 Geofence
Common issues encountered with geofences are listed in this section.
[Link] Name and Description Fields are not displayed for the DB-based Geofences
If name and description fields are not displayed for database-based geofence, ensure to follow
steps mentioned below:
1. Go to Catalog and click Edit for the required database-based geofence.
2. Click Edit for Source Type Properties and then Next.
3. Ensure that the mapping for Name and Description is defined in Shape section.
4. Once these mappings are defined, you can see the name and description for the geofence.
10.2.7 Cube
Common issues encountered with cubes are listed in this section.
10-17
Chapter 10
Common Issues and Remedies
10.2.8 Dashboard
Common issues encountered with dashboards are listed in this section.
10-18
Chapter 10
Common Issues and Remedies
Ensure that CQL Queries for Each Query Stage Emit Output
To check if the CQL queries are emitting output events to monitor CQL Queries using CQL
Engine Metrics:
1. Open CQL Engine Query Details page. For more information, see Access CQL Engine
Metrics.
2. Check that at least one partition has Total Output Events greater than zero under the
Execution Statistics section.
If your query is running without any error and input data is continuously coming, then the
Total Output Events will keep rising.
10-19
Chapter 10
Common Issues and Remedies
6. Listen to the Kafka topic where output of the stage is being pushed.
Since this topic is created using Kafka APIs, you cannot consume this topic with REST
APIs. Follow these steps to listen to the Kafka topic:
a. Listen to the Kafka topic hosted on a standard Apache Kafka installation.
You can listen to the Kafka topic using utilities from a Kafka Installation. kafka-
[Link] is a utility script available as part of any Kafka installation.
To listen to Kafka topic:
i. Determine the Zookeeper Address from Apache Kafka Installation based Cluster.
ii. Use following command to listen the Kafka topic:
10-20
Chapter 10
Common Issues and Remedies
This exception usually occurs when you do not have free resources on your cluster.
Workaround:
Use external Spark cluster or get better machine and configure the cluster with more
resources.
10-21