Latest commit

History

33 Commits

Folders and files

NameName
Last commit message
Last commit date

Repository files navigation

Stream Processing using KSQL

This project show how to use KSQL (Streaming SQL Engine for Apache Kafka) to stream processing.

Enviroment

For you to use this repository you will need the following softwares:

However, only Docker and Docker Compose need is installed in your machine. All Kafka ecosystem will be embedded via docker images.

Steps

  1. Install Python and Pip
  2. Install Docker and Docker Compose
  3. Load Images
  4. Create Topics
  5. Start Simulator

1 - Install Python and Pip

The installation process of the Python and Pip is very easy. So this tutorial dont't will cover this steps. I recommend you look for more information in www.python.org and pip.pypa.io.

After you install Python and Pip run the command below to install all dependencies need to execute click_simulator.py application. This code is responsible to simulate the click events into an web page. It will generate unbounded click events, sending a flow continuous messages to a Kafka topic.

pip install -r requirements.txt

2 - Install Docker and Docker Compose

This tutorial does not demonstrate the installation process for Docker and Docker Compose. I strongly recommend you to visit the Docker installation link for more informations. Please click here.

3 - Loading Images

docker-compose up

or

docker-compose up -d

The last command allow you to run docker-compose in the background.

4 - Create Topics

docker-compose exec kafka kafka-topics --create --topic com.mywebsite.streams.pages --bootstrap-server localhost:9092
docker-compose exec kafka kafka-topics --create --topic com.mywebsite.streams.clickevents --bootstrap-server localhost:9092

5 - Start Simulator

python click_simulator.py

If you executed all steps correctly. You will see an image similar that below.

Starting application
Message: {"email": "anoble@yahoo.com", "timestamp": "1986-03-10T16:38:40", "uri": "https://mitchell.info/login.php", "number": 358}
Message: {"email": "leonardpatrick@mason-clark.info", "timestamp": "1971-04-25T10:09:26", "uri": "https://www.bailey.com/search/about/", "number": 431}
Message: {"email": "morriskatie@villarreal-villa.biz", "timestamp": "1996-11-22T00:12:20", "uri": "http://www.woodard.info/terms.php", "number": 838}
Message: {"email": "kenneth79@rogers.info", "timestamp": "2005-10-24T22:16:59", "uri": "http://www.king.com/wp-content/blog/blog/index/", "number": 793}
Message: {"email": "wbailey@wu-martinez.net", "timestamp": "1995-06-20T12:44:44", "uri": "https://www.smith-neal.com/categories/login/", "number": 509}
Message: {"email": "tkennedy@hall-wolfe.org", "timestamp": "2009-01-27T14:04:20", "uri": "https://www.marshall-holmes.info/", "number": 336}
Message: {"email": "steven15@yahoo.com", "timestamp": "2019-12-13T16:09:11", "uri": "https://www.sims.net/main.html", "number": 263}
Message: {"email": "hobbsmario@hotmail.com", "timestamp": "1990-08-16T05:09:04", "uri": "http://www.smith.com/search/tags/explore/about.jsp", "number": 61}
...

Connecting to KSQL Server

docker-compose exec ksql ksql http://localhost:8088

After you connect to KSQL Server you will see the image below:

 ===========================================
= _ __ _____ ____ _ =
= ||/ // ____|/ __ \|| =
= |' /| (___ | | | | | = = | < \___ \| | | | | = = | . \ ____) | |__| | |____ = = |_|\_\_____/ \___\_\______| = = = = Streaming SQL Engine for Apache Kafka® = ===========================================Copyright 2017-2019 Confluent Inc.CLI v5.4.1, Server v5.4.1 located at http://localhost:8088Having trouble? Type 'help' (case-insensitive) for a rundown of how things work!ksql>

Some Commands

Show all topics

ksql> SHOW TOPICS;
Kafka Topic | Partitions | Partition Replicas
---------------------------------------------------------------------
com.mywebsite.streams.clickevents | 5 | 1
com.mywebsite.streams.pages | 1 | 1
---------------------------------------------------------------------

Show all streams

ksql> SHOW STREAMS;
Stream Name | Kafka Topic | Format
--------------------------------------------------------
CLICKEVENTS | com.mywebsite.streams.clickevents | JSON
--------------------------------------------------------

Creating a Stream

If you need run it in the background mode.

CREATE STREAM clickevents
(email VARCHAR,
timestampVARCHAR,
uri VARCHAR,
numberINTEGER)
WITH (KAFKA_TOPIC='com.mywebsite.streams.clickevents',
VALUE_FORMAT='JSON');

Creating a Table

CREATETABLEpages
(uri VARCHAR,
description VARCHAR,
created VARCHAR)
WITH (KAFKA_TOPIC='com.mywebsite.streams.pages',
VALUE_FORMAT='JSON',
KEY='uri');

Creating a Table from a Query

CREATETABLEa_pagesASSELECT*FROM pages WHERE uri LIKE'http://www.a%';

Querying a Table or Stream

SELECT*FROM clickevents EMIT CHANGES;

Describing a Table and Stream

ksql> DESCRIBE PAGES;
Name : PAGES
Field | Type
-----------------------------------------
ROWTIME | BIGINT (system)
ROWKEY | VARCHAR(STRING) (system)
URI | VARCHAR(STRING)
DESCRIPTION | VARCHAR(STRING)
CREATED | VARCHAR(STRING)
-----------------------------------------
For runtime statistics and query details run: DESCRIBE EXTENDED <Stream,Table>;

Managing Offsets

Like all Kafka Consumers, KSQL by default begins consumption at the latest offset. This can be a problem for some scenarios. In the following example we're going to create a pages table -- but -- we want all the data available to us in this table. In other words, we want KSQL to start from the earliest offset. To do this, we will use the SET command to set the configuration variabl auto.offset.reset for our session -- and before we run any commands.

SET 'auto.offset.reset' = 'earliest';

Also note that this can be set at the KSQL server level, if you'd like. Once you're done querying or creating tables or streams with this value, you can set it back to its original setting by simply running:

UNSET 'auto.offset.reset';

Scalar Functions

KSQL Provides a number of Scalar functions for us to make use of.

Lets write a function that takes advantage of some of these features:

SELECT UCASE(SUBSTRING(uri, 12))
FROM clickevents
WHERE number > 100
AND uri LIKE 'http://www.k%' EMIT CHANGES;

Notice that as soon as you hit CTRL+C your query ends

Deleting a Table

As with Streams, we must first find the running underlying query, and then drop the table. First, find your query:

ksql> SHOW QUERIES;
Query ID | Kafka Topic | Query String
----------------------------------------------------------------------------------------------
CTAS_A_PAGES_1 | A_PAGES | CREATE TABLE a_pages AS
SELECT * FROM pages WHERE uri LIKE 'http://www.a%';
----------------------------------------------------------------------------------------------
For detailed information on a Query run: EXPLAIN <Query ID>;

Find your query, which in this case is CTAS_A_PAGES_1 and then, finally, TERMINATE the query and DROP the table:

TERMINATE QUERY CTAS_A_PAGES_1;
DROP TABLE A_PAGES;

Windowing

Hopping and Tumbling Windows

In this demonstration we'll see how to create Tables with windowing enabled.

Tumbling Windows

Let's create a tumbling clickevents table, where the window size is 30 seconds.

CREATE STREAM clickevents_tumbling ASSELECT*FROM clickevents
WINDOW TUMBLING (SIZE 30 SECONDS);

Hopping Windows

Now we can create a Table with a hopping window of 30 seconds with 5 second increments.

CREATETABLEclickevents_hoppingASSELECT uri FROM clickevents
WINDOW HOPPING (SIZE 30 SECONDS, ADVANCE BY 5 SECONDS)
WHERE uri LIKE'http://www.b%'GROUP BY uri;

The above window is 30 seconds long and advances by 5 second. If you query the table you will see the associated window times!

Session Windows

Finally, lets see how session windows work. We're going to define the session as 5 minutes in order to group many events to the same window

CREATETABLEclickevents_sessionASSELECT uri FROM clickevents
WINDOW SESSION (5 MINUTES)
WHERE uri LIKE'http://www.b%'GROUP BY uri;

Kafka CLI Basic Commands

Creating a topic:

docker-compose exec kafka kafka-topics --create --topic <topic-name> --bootstrap-server localhost:9092

Writing a topic:

docker-compose exec kafka kafka-console-producer --topic <topic-name> --bootstrap-server localhost:9092

You must type on console the press key enter.

Reading a topic:

docker-compose exec kafka kafka-console-consumer.sh --topic <topic-name> --from-beginning --bootstrap-server localhost:9092

Contributing

Pull requests are welcome. For major changes, please open an issue first to discuss what you would like to change.

Please make sure to update tests as appropriate.

License

MIT

About

This project show how to use KSQL (Streaming SQL Engine for Apache Kafka) to stream processing.

Topics

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content

Latest commit

History

33 Commits

Folders and files

NameName
Last commit message
Last commit date

Repository files navigation

Stream Processing using KSQL

This project show how to use KSQL (Streaming SQL Engine for Apache Kafka) to stream processing.

Enviroment

For you to use this repository you will need the following softwares:

However, only Docker and Docker Compose need is installed in your machine. All Kafka ecosystem will be embedded via docker images.

Steps

  1. Install Python and Pip
  2. Install Docker and Docker Compose
  3. Load Images
  4. Create Topics
  5. Start Simulator

1 - Install Python and Pip

The installation process of the Python and Pip is very easy. So this tutorial dont't will cover this steps. I recommend you look for more information in www.python.org and pip.pypa.io.

After you install Python and Pip run the command below to install all dependencies need to execute click_simulator.py application. This code is responsible to simulate the click events into an web page. It will generate unbounded click events, sending a flow continuous messages to a Kafka topic.

pip install -r requirements.txt

2 - Install Docker and Docker Compose

This tutorial does not demonstrate the installation process for Docker and Docker Compose. I strongly recommend you to visit the Docker installation link for more informations. Please click here.

3 - Loading Images

docker-compose up

or

docker-compose up -d

The last command allow you to run docker-compose in the background.

4 - Create Topics

docker-compose exec kafka kafka-topics --create --topic com.mywebsite.streams.pages --bootstrap-server localhost:9092
docker-compose exec kafka kafka-topics --create --topic com.mywebsite.streams.clickevents --bootstrap-server localhost:9092

5 - Start Simulator

python click_simulator.py

If you executed all steps correctly. You will see an image similar that below.

Starting application
Message: {"email": "anoble@yahoo.com", "timestamp": "1986-03-10T16:38:40", "uri": "https://mitchell.info/login.php", "number": 358}
Message: {"email": "leonardpatrick@mason-clark.info", "timestamp": "1971-04-25T10:09:26", "uri": "https://www.bailey.com/search/about/", "number": 431}
Message: {"email": "morriskatie@villarreal-villa.biz", "timestamp": "1996-11-22T00:12:20", "uri": "http://www.woodard.info/terms.php", "number": 838}
Message: {"email": "kenneth79@rogers.info", "timestamp": "2005-10-24T22:16:59", "uri": "http://www.king.com/wp-content/blog/blog/index/", "number": 793}
Message: {"email": "wbailey@wu-martinez.net", "timestamp": "1995-06-20T12:44:44", "uri": "https://www.smith-neal.com/categories/login/", "number": 509}
Message: {"email": "tkennedy@hall-wolfe.org", "timestamp": "2009-01-27T14:04:20", "uri": "https://www.marshall-holmes.info/", "number": 336}
Message: {"email": "steven15@yahoo.com", "timestamp": "2019-12-13T16:09:11", "uri": "https://www.sims.net/main.html", "number": 263}
Message: {"email": "hobbsmario@hotmail.com", "timestamp": "1990-08-16T05:09:04", "uri": "http://www.smith.com/search/tags/explore/about.jsp", "number": 61}
...

Connecting to KSQL Server

docker-compose exec ksql ksql http://localhost:8088

After you connect to KSQL Server you will see the image below:

 ===========================================
= _ __ _____ ____ _ =
= ||/ // ____|/ __ \|| =
= |' /| (___ | | | | | = = | < \___ \| | | | | = = | . \ ____) | |__| | |____ = = |_|\_\_____/ \___\_\______| = = = = Streaming SQL Engine for Apache Kafka® = ===========================================Copyright 2017-2019 Confluent Inc.CLI v5.4.1, Server v5.4.1 located at http://localhost:8088Having trouble? Type 'help' (case-insensitive) for a rundown of how things work!ksql>

Some Commands

Show all topics

ksql> SHOW TOPICS;
Kafka Topic | Partitions | Partition Replicas
---------------------------------------------------------------------
com.mywebsite.streams.clickevents | 5 | 1
com.mywebsite.streams.pages | 1 | 1
---------------------------------------------------------------------

Show all streams

ksql> SHOW STREAMS;
Stream Name | Kafka Topic | Format
--------------------------------------------------------
CLICKEVENTS | com.mywebsite.streams.clickevents | JSON
--------------------------------------------------------

Creating a Stream

If you need run it in the background mode.

CREATE STREAM clickevents
(email VARCHAR,
timestampVARCHAR,
uri VARCHAR,
numberINTEGER)
WITH (KAFKA_TOPIC='com.mywebsite.streams.clickevents',
VALUE_FORMAT='JSON');

Creating a Table

CREATETABLEpages
(uri VARCHAR,
description VARCHAR,
created VARCHAR)
WITH (KAFKA_TOPIC='com.mywebsite.streams.pages',
VALUE_FORMAT='JSON',
KEY='uri');

Creating a Table from a Query

CREATETABLEa_pagesASSELECT*FROM pages WHERE uri LIKE'http://www.a%';

Querying a Table or Stream

SELECT*FROM clickevents EMIT CHANGES;

Describing a Table and Stream

ksql> DESCRIBE PAGES;
Name : PAGES
Field | Type
-----------------------------------------
ROWTIME | BIGINT (system)
ROWKEY | VARCHAR(STRING) (system)
URI | VARCHAR(STRING)
DESCRIPTION | VARCHAR(STRING)
CREATED | VARCHAR(STRING)
-----------------------------------------
For runtime statistics and query details run: DESCRIBE EXTENDED <Stream,Table>;

Managing Offsets

Like all Kafka Consumers, KSQL by default begins consumption at the latest offset. This can be a problem for some scenarios. In the following example we're going to create a pages table -- but -- we want all the data available to us in this table. In other words, we want KSQL to start from the earliest offset. To do this, we will use the SET command to set the configuration variabl auto.offset.reset for our session -- and before we run any commands.

SET 'auto.offset.reset' = 'earliest';

Also note that this can be set at the KSQL server level, if you'd like. Once you're done querying or creating tables or streams with this value, you can set it back to its original setting by simply running:

UNSET 'auto.offset.reset';

Scalar Functions

KSQL Provides a number of Scalar functions for us to make use of.

Lets write a function that takes advantage of some of these features:

SELECT UCASE(SUBSTRING(uri, 12))
FROM clickevents
WHERE number > 100
AND uri LIKE 'http://www.k%' EMIT CHANGES;

Notice that as soon as you hit CTRL+C your query ends

Deleting a Table

As with Streams, we must first find the running underlying query, and then drop the table. First, find your query:

ksql> SHOW QUERIES;
Query ID | Kafka Topic | Query String
----------------------------------------------------------------------------------------------
CTAS_A_PAGES_1 | A_PAGES | CREATE TABLE a_pages AS
SELECT * FROM pages WHERE uri LIKE 'http://www.a%';
----------------------------------------------------------------------------------------------
For detailed information on a Query run: EXPLAIN <Query ID>;

Find your query, which in this case is CTAS_A_PAGES_1 and then, finally, TERMINATE the query and DROP the table:

TERMINATE QUERY CTAS_A_PAGES_1;
DROP TABLE A_PAGES;

Windowing

Hopping and Tumbling Windows

In this demonstration we'll see how to create Tables with windowing enabled.

Tumbling Windows

Let's create a tumbling clickevents table, where the window size is 30 seconds.

CREATE STREAM clickevents_tumbling ASSELECT*FROM clickevents
WINDOW TUMBLING (SIZE 30 SECONDS);

Hopping Windows

Now we can create a Table with a hopping window of 30 seconds with 5 second increments.

CREATETABLEclickevents_hoppingASSELECT uri FROM clickevents
WINDOW HOPPING (SIZE 30 SECONDS, ADVANCE BY 5 SECONDS)
WHERE uri LIKE'http://www.b%'GROUP BY uri;

The above window is 30 seconds long and advances by 5 second. If you query the table you will see the associated window times!

Session Windows

Finally, lets see how session windows work. We're going to define the session as 5 minutes in order to group many events to the same window

CREATETABLEclickevents_sessionASSELECT uri FROM clickevents
WINDOW SESSION (5 MINUTES)
WHERE uri LIKE'http://www.b%'GROUP BY uri;

Kafka CLI Basic Commands

Creating a topic:

docker-compose exec kafka kafka-topics --create --topic <topic-name> --bootstrap-server localhost:9092

Writing a topic:

docker-compose exec kafka kafka-console-producer --topic <topic-name> --bootstrap-server localhost:9092

You must type on console the press key enter.

Reading a topic:

docker-compose exec kafka kafka-console-consumer.sh --topic <topic-name> --from-beginning --bootstrap-server localhost:9092

Contributing

Pull requests are welcome. For major changes, please open an issue first to discuss what you would like to change.

Please make sure to update tests as appropriate.

License

MIT

About

This project show how to use KSQL (Streaming SQL Engine for Apache Kafka) to stream processing.

Topics

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Latest commit

History

33 Commits

Folders and files

NameName
Last commit message
Last commit date

Repository files navigation

Stream Processing using KSQL

This project show how to use KSQL (Streaming SQL Engine for Apache Kafka) to stream processing.

Enviroment

For you to use this repository you will need the following softwares:

However, only Docker and Docker Compose need is installed in your machine. All Kafka ecosystem will be embedded via docker images.

Steps

  1. Install Python and Pip
  2. Install Docker and Docker Compose
  3. Load Images
  4. Create Topics
  5. Start Simulator

1 - Install Python and Pip

The installation process of the Python and Pip is very easy. So this tutorial dont't will cover this steps. I recommend you look for more information in www.python.org and pip.pypa.io.

After you install Python and Pip run the command below to install all dependencies need to execute click_simulator.py application. This code is responsible to simulate the click events into an web page. It will generate unbounded click events, sending a flow continuous messages to a Kafka topic.

pip install -r requirements.txt

2 - Install Docker and Docker Compose

This tutorial does not demonstrate the installation process for Docker and Docker Compose. I strongly recommend you to visit the Docker installation link for more informations. Please click here.

3 - Loading Images

docker-compose up

or

docker-compose up -d

The last command allow you to run docker-compose in the background.

4 - Create Topics

docker-compose exec kafka kafka-topics --create --topic com.mywebsite.streams.pages --bootstrap-server localhost:9092
docker-compose exec kafka kafka-topics --create --topic com.mywebsite.streams.clickevents --bootstrap-server localhost:9092

5 - Start Simulator

python click_simulator.py

If you executed all steps correctly. You will see an image similar that below.

Starting application
Message: {"email": "anoble@yahoo.com", "timestamp": "1986-03-10T16:38:40", "uri": "https://mitchell.info/login.php", "number": 358}
Message: {"email": "leonardpatrick@mason-clark.info", "timestamp": "1971-04-25T10:09:26", "uri": "https://www.bailey.com/search/about/", "number": 431}
Message: {"email": "morriskatie@villarreal-villa.biz", "timestamp": "1996-11-22T00:12:20", "uri": "http://www.woodard.info/terms.php", "number": 838}
Message: {"email": "kenneth79@rogers.info", "timestamp": "2005-10-24T22:16:59", "uri": "http://www.king.com/wp-content/blog/blog/index/", "number": 793}
Message: {"email": "wbailey@wu-martinez.net", "timestamp": "1995-06-20T12:44:44", "uri": "https://www.smith-neal.com/categories/login/", "number": 509}
Message: {"email": "tkennedy@hall-wolfe.org", "timestamp": "2009-01-27T14:04:20", "uri": "https://www.marshall-holmes.info/", "number": 336}
Message: {"email": "steven15@yahoo.com", "timestamp": "2019-12-13T16:09:11", "uri": "https://www.sims.net/main.html", "number": 263}
Message: {"email": "hobbsmario@hotmail.com", "timestamp": "1990-08-16T05:09:04", "uri": "http://www.smith.com/search/tags/explore/about.jsp", "number": 61}
...

Connecting to KSQL Server

docker-compose exec ksql ksql http://localhost:8088

After you connect to KSQL Server you will see the image below:

 ===========================================
= _ __ _____ ____ _ =
= ||/ // ____|/ __ \|| =
= |' /| (___ | | | | | = = | < \___ \| | | | | = = | . \ ____) | |__| | |____ = = |_|\_\_____/ \___\_\______| = = = = Streaming SQL Engine for Apache Kafka® = ===========================================Copyright 2017-2019 Confluent Inc.CLI v5.4.1, Server v5.4.1 located at http://localhost:8088Having trouble? Type 'help' (case-insensitive) for a rundown of how things work!ksql>

Some Commands

Show all topics

ksql> SHOW TOPICS;
Kafka Topic | Partitions | Partition Replicas
---------------------------------------------------------------------
com.mywebsite.streams.clickevents | 5 | 1
com.mywebsite.streams.pages | 1 | 1
---------------------------------------------------------------------

Show all streams

ksql> SHOW STREAMS;
Stream Name | Kafka Topic | Format
--------------------------------------------------------
CLICKEVENTS | com.mywebsite.streams.clickevents | JSON
--------------------------------------------------------

Creating a Stream

If you need run it in the background mode.

CREATE STREAM clickevents
(email VARCHAR,
timestampVARCHAR,
uri VARCHAR,
numberINTEGER)
WITH (KAFKA_TOPIC='com.mywebsite.streams.clickevents',
VALUE_FORMAT='JSON');

Creating a Table

CREATETABLEpages
(uri VARCHAR,
description VARCHAR,
created VARCHAR)
WITH (KAFKA_TOPIC='com.mywebsite.streams.pages',
VALUE_FORMAT='JSON',
KEY='uri');

Creating a Table from a Query

CREATETABLEa_pagesASSELECT*FROM pages WHERE uri LIKE'http://www.a%';

Querying a Table or Stream

SELECT*FROM clickevents EMIT CHANGES;

Describing a Table and Stream

ksql> DESCRIBE PAGES;
Name : PAGES
Field | Type
-----------------------------------------
ROWTIME | BIGINT (system)
ROWKEY | VARCHAR(STRING) (system)
URI | VARCHAR(STRING)
DESCRIPTION | VARCHAR(STRING)
CREATED | VARCHAR(STRING)
-----------------------------------------
For runtime statistics and query details run: DESCRIBE EXTENDED <Stream,Table>;

Managing Offsets

Like all Kafka Consumers, KSQL by default begins consumption at the latest offset. This can be a problem for some scenarios. In the following example we're going to create a pages table -- but -- we want all the data available to us in this table. In other words, we want KSQL to start from the earliest offset. To do this, we will use the SET command to set the configuration variabl auto.offset.reset for our session -- and before we run any commands.

SET 'auto.offset.reset' = 'earliest';

Also note that this can be set at the KSQL server level, if you'd like. Once you're done querying or creating tables or streams with this value, you can set it back to its original setting by simply running:

UNSET 'auto.offset.reset';

Scalar Functions

KSQL Provides a number of Scalar functions for us to make use of.

Lets write a function that takes advantage of some of these features:

SELECT UCASE(SUBSTRING(uri, 12))
FROM clickevents
WHERE number > 100
AND uri LIKE 'http://www.k%' EMIT CHANGES;

Notice that as soon as you hit CTRL+C your query ends

Deleting a Table

As with Streams, we must first find the running underlying query, and then drop the table. First, find your query:

ksql> SHOW QUERIES;
Query ID | Kafka Topic | Query String
----------------------------------------------------------------------------------------------
CTAS_A_PAGES_1 | A_PAGES | CREATE TABLE a_pages AS
SELECT * FROM pages WHERE uri LIKE 'http://www.a%';
----------------------------------------------------------------------------------------------
For detailed information on a Query run: EXPLAIN <Query ID>;

Find your query, which in this case is CTAS_A_PAGES_1 and then, finally, TERMINATE the query and DROP the table:

TERMINATE QUERY CTAS_A_PAGES_1;
DROP TABLE A_PAGES;

Windowing

Hopping and Tumbling Windows

In this demonstration we'll see how to create Tables with windowing enabled.

Tumbling Windows

Let's create a tumbling clickevents table, where the window size is 30 seconds.

CREATE STREAM clickevents_tumbling ASSELECT*FROM clickevents
WINDOW TUMBLING (SIZE 30 SECONDS);

Hopping Windows

Now we can create a Table with a hopping window of 30 seconds with 5 second increments.

CREATETABLEclickevents_hoppingASSELECT uri FROM clickevents
WINDOW HOPPING (SIZE 30 SECONDS, ADVANCE BY 5 SECONDS)
WHERE uri LIKE'http://www.b%'GROUP BY uri;

The above window is 30 seconds long and advances by 5 second. If you query the table you will see the associated window times!

Session Windows

Finally, lets see how session windows work. We're going to define the session as 5 minutes in order to group many events to the same window

CREATETABLEclickevents_sessionASSELECT uri FROM clickevents
WINDOW SESSION (5 MINUTES)
WHERE uri LIKE'http://www.b%'GROUP BY uri;

Kafka CLI Basic Commands

Creating a topic:

docker-compose exec kafka kafka-topics --create --topic <topic-name> --bootstrap-server localhost:9092

Writing a topic:

docker-compose exec kafka kafka-console-producer --topic <topic-name> --bootstrap-server localhost:9092

You must type on console the press key enter.

Reading a topic:

docker-compose exec kafka kafka-console-consumer.sh --topic <topic-name> --from-beginning --bootstrap-server localhost:9092

Contributing

Pull requests are welcome. For major changes, please open an issue first to discuss what you would like to change.

Please make sure to update tests as appropriate.

License

MIT

About

This project show how to use KSQL (Streaming SQL Engine for Apache Kafka) to stream processing.

Topics

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Latest commit

History

33 Commits

Folders and files

NameName
Last commit message
Last commit date

Repository files navigation

Stream Processing using KSQL

This project show how to use KSQL (Streaming SQL Engine for Apache Kafka) to stream processing.

Enviroment

For you to use this repository you will need the following softwares:

However, only Docker and Docker Compose need is installed in your machine. All Kafka ecosystem will be embedded via docker images.

Steps

  1. Install Python and Pip
  2. Install Docker and Docker Compose
  3. Load Images
  4. Create Topics
  5. Start Simulator

1 - Install Python and Pip

The installation process of the Python and Pip is very easy. So this tutorial dont't will cover this steps. I recommend you look for more information in www.python.org and pip.pypa.io.

After you install Python and Pip run the command below to install all dependencies need to execute click_simulator.py application. This code is responsible to simulate the click events into an web page. It will generate unbounded click events, sending a flow continuous messages to a Kafka topic.

pip install -r requirements.txt

2 - Install Docker and Docker Compose

This tutorial does not demonstrate the installation process for Docker and Docker Compose. I strongly recommend you to visit the Docker installation link for more informations. Please click here.

3 - Loading Images

docker-compose up

or

docker-compose up -d

The last command allow you to run docker-compose in the background.

4 - Create Topics

docker-compose exec kafka kafka-topics --create --topic com.mywebsite.streams.pages --bootstrap-server localhost:9092
docker-compose exec kafka kafka-topics --create --topic com.mywebsite.streams.clickevents --bootstrap-server localhost:9092

5 - Start Simulator

python click_simulator.py

If you executed all steps correctly. You will see an image similar that below.

Starting application
Message: {"email": "anoble@yahoo.com", "timestamp": "1986-03-10T16:38:40", "uri": "https://mitchell.info/login.php", "number": 358}
Message: {"email": "leonardpatrick@mason-clark.info", "timestamp": "1971-04-25T10:09:26", "uri": "https://www.bailey.com/search/about/", "number": 431}
Message: {"email": "morriskatie@villarreal-villa.biz", "timestamp": "1996-11-22T00:12:20", "uri": "http://www.woodard.info/terms.php", "number": 838}
Message: {"email": "kenneth79@rogers.info", "timestamp": "2005-10-24T22:16:59", "uri": "http://www.king.com/wp-content/blog/blog/index/", "number": 793}
Message: {"email": "wbailey@wu-martinez.net", "timestamp": "1995-06-20T12:44:44", "uri": "https://www.smith-neal.com/categories/login/", "number": 509}
Message: {"email": "tkennedy@hall-wolfe.org", "timestamp": "2009-01-27T14:04:20", "uri": "https://www.marshall-holmes.info/", "number": 336}
Message: {"email": "steven15@yahoo.com", "timestamp": "2019-12-13T16:09:11", "uri": "https://www.sims.net/main.html", "number": 263}
Message: {"email": "hobbsmario@hotmail.com", "timestamp": "1990-08-16T05:09:04", "uri": "http://www.smith.com/search/tags/explore/about.jsp", "number": 61}
...

Connecting to KSQL Server

docker-compose exec ksql ksql http://localhost:8088

After you connect to KSQL Server you will see the image below:

 ===========================================
= _ __ _____ ____ _ =
= ||/ // ____|/ __ \|| =
= |' /| (___ | | | | | = = | < \___ \| | | | | = = | . \ ____) | |__| | |____ = = |_|\_\_____/ \___\_\______| = = = = Streaming SQL Engine for Apache Kafka® = ===========================================Copyright 2017-2019 Confluent Inc.CLI v5.4.1, Server v5.4.1 located at http://localhost:8088Having trouble? Type 'help' (case-insensitive) for a rundown of how things work!ksql>

Some Commands

Show all topics

ksql> SHOW TOPICS;
Kafka Topic | Partitions | Partition Replicas
---------------------------------------------------------------------
com.mywebsite.streams.clickevents | 5 | 1
com.mywebsite.streams.pages | 1 | 1
---------------------------------------------------------------------

Show all streams

ksql> SHOW STREAMS;
Stream Name | Kafka Topic | Format
--------------------------------------------------------
CLICKEVENTS | com.mywebsite.streams.clickevents | JSON
--------------------------------------------------------

Creating a Stream

If you need run it in the background mode.

CREATE STREAM clickevents
(email VARCHAR,
timestampVARCHAR,
uri VARCHAR,
numberINTEGER)
WITH (KAFKA_TOPIC='com.mywebsite.streams.clickevents',
VALUE_FORMAT='JSON');

Creating a Table

CREATETABLEpages
(uri VARCHAR,
description VARCHAR,
created VARCHAR)
WITH (KAFKA_TOPIC='com.mywebsite.streams.pages',
VALUE_FORMAT='JSON',
KEY='uri');

Creating a Table from a Query

CREATETABLEa_pagesASSELECT*FROM pages WHERE uri LIKE'http://www.a%';

Querying a Table or Stream

SELECT*FROM clickevents EMIT CHANGES;

Describing a Table and Stream

ksql> DESCRIBE PAGES;
Name : PAGES
Field | Type
-----------------------------------------
ROWTIME | BIGINT (system)
ROWKEY | VARCHAR(STRING) (system)
URI | VARCHAR(STRING)
DESCRIPTION | VARCHAR(STRING)
CREATED | VARCHAR(STRING)
-----------------------------------------
For runtime statistics and query details run: DESCRIBE EXTENDED <Stream,Table>;

Managing Offsets

Like all Kafka Consumers, KSQL by default begins consumption at the latest offset. This can be a problem for some scenarios. In the following example we're going to create a pages table -- but -- we want all the data available to us in this table. In other words, we want KSQL to start from the earliest offset. To do this, we will use the SET command to set the configuration variabl auto.offset.reset for our session -- and before we run any commands.

SET 'auto.offset.reset' = 'earliest';

Also note that this can be set at the KSQL server level, if you'd like. Once you're done querying or creating tables or streams with this value, you can set it back to its original setting by simply running:

UNSET 'auto.offset.reset';

Scalar Functions

KSQL Provides a number of Scalar functions for us to make use of.

Lets write a function that takes advantage of some of these features:

SELECT UCASE(SUBSTRING(uri, 12))
FROM clickevents
WHERE number > 100
AND uri LIKE 'http://www.k%' EMIT CHANGES;

Notice that as soon as you hit CTRL+C your query ends

Deleting a Table

As with Streams, we must first find the running underlying query, and then drop the table. First, find your query:

ksql> SHOW QUERIES;
Query ID | Kafka Topic | Query String
----------------------------------------------------------------------------------------------
CTAS_A_PAGES_1 | A_PAGES | CREATE TABLE a_pages AS
SELECT * FROM pages WHERE uri LIKE 'http://www.a%';
----------------------------------------------------------------------------------------------
For detailed information on a Query run: EXPLAIN <Query ID>;

Find your query, which in this case is CTAS_A_PAGES_1 and then, finally, TERMINATE the query and DROP the table:

TERMINATE QUERY CTAS_A_PAGES_1;
DROP TABLE A_PAGES;

Windowing

Hopping and Tumbling Windows

In this demonstration we'll see how to create Tables with windowing enabled.

Tumbling Windows

Let's create a tumbling clickevents table, where the window size is 30 seconds.

CREATE STREAM clickevents_tumbling ASSELECT*FROM clickevents
WINDOW TUMBLING (SIZE 30 SECONDS);

Hopping Windows

Now we can create a Table with a hopping window of 30 seconds with 5 second increments.

CREATETABLEclickevents_hoppingASSELECT uri FROM clickevents
WINDOW HOPPING (SIZE 30 SECONDS, ADVANCE BY 5 SECONDS)
WHERE uri LIKE'http://www.b%'GROUP BY uri;

The above window is 30 seconds long and advances by 5 second. If you query the table you will see the associated window times!

Session Windows

Finally, lets see how session windows work. We're going to define the session as 5 minutes in order to group many events to the same window

CREATETABLEclickevents_sessionASSELECT uri FROM clickevents
WINDOW SESSION (5 MINUTES)
WHERE uri LIKE'http://www.b%'GROUP BY uri;

Kafka CLI Basic Commands

Creating a topic:

docker-compose exec kafka kafka-topics --create --topic <topic-name> --bootstrap-server localhost:9092

Writing a topic:

docker-compose exec kafka kafka-console-producer --topic <topic-name> --bootstrap-server localhost:9092

You must type on console the press key enter.

Reading a topic:

docker-compose exec kafka kafka-console-consumer.sh --topic <topic-name> --from-beginning --bootstrap-server localhost:9092

Contributing

Pull requests are welcome. For major changes, please open an issue first to discuss what you would like to change.

Please make sure to update tests as appropriate.

License

MIT

About

This project show how to use KSQL (Streaming SQL Engine for Apache Kafka) to stream processing.

Topics

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content

Latest commit

History

33 Commits

Folders and files

NameName
Last commit message
Last commit date

Repository files navigation

Stream Processing using KSQL

This project show how to use KSQL (Streaming SQL Engine for Apache Kafka) to stream processing.

Enviroment

For you to use this repository you will need the following softwares:

However, only Docker and Docker Compose need is installed in your machine. All Kafka ecosystem will be embedded via docker images.

Steps

  1. Install Python and Pip
  2. Install Docker and Docker Compose
  3. Load Images
  4. Create Topics
  5. Start Simulator

1 - Install Python and Pip

The installation process of the Python and Pip is very easy. So this tutorial dont't will cover this steps. I recommend you look for more information in www.python.org and pip.pypa.io.

After you install Python and Pip run the command below to install all dependencies need to execute click_simulator.py application. This code is responsible to simulate the click events into an web page. It will generate unbounded click events, sending a flow continuous messages to a Kafka topic.

pip install -r requirements.txt

2 - Install Docker and Docker Compose

This tutorial does not demonstrate the installation process for Docker and Docker Compose. I strongly recommend you to visit the Docker installation link for more informations. Please click here.

3 - Loading Images

docker-compose up

or

docker-compose up -d

The last command allow you to run docker-compose in the background.

4 - Create Topics

docker-compose exec kafka kafka-topics --create --topic com.mywebsite.streams.pages --bootstrap-server localhost:9092
docker-compose exec kafka kafka-topics --create --topic com.mywebsite.streams.clickevents --bootstrap-server localhost:9092

5 - Start Simulator

python click_simulator.py

If you executed all steps correctly. You will see an image similar that below.

Starting application
Message: {"email": "anoble@yahoo.com", "timestamp": "1986-03-10T16:38:40", "uri": "https://mitchell.info/login.php", "number": 358}
Message: {"email": "leonardpatrick@mason-clark.info", "timestamp": "1971-04-25T10:09:26", "uri": "https://www.bailey.com/search/about/", "number": 431}
Message: {"email": "morriskatie@villarreal-villa.biz", "timestamp": "1996-11-22T00:12:20", "uri": "http://www.woodard.info/terms.php", "number": 838}
Message: {"email": "kenneth79@rogers.info", "timestamp": "2005-10-24T22:16:59", "uri": "http://www.king.com/wp-content/blog/blog/index/", "number": 793}
Message: {"email": "wbailey@wu-martinez.net", "timestamp": "1995-06-20T12:44:44", "uri": "https://www.smith-neal.com/categories/login/", "number": 509}
Message: {"email": "tkennedy@hall-wolfe.org", "timestamp": "2009-01-27T14:04:20", "uri": "https://www.marshall-holmes.info/", "number": 336}
Message: {"email": "steven15@yahoo.com", "timestamp": "2019-12-13T16:09:11", "uri": "https://www.sims.net/main.html", "number": 263}
Message: {"email": "hobbsmario@hotmail.com", "timestamp": "1990-08-16T05:09:04", "uri": "http://www.smith.com/search/tags/explore/about.jsp", "number": 61}
...

Connecting to KSQL Server

docker-compose exec ksql ksql http://localhost:8088

After you connect to KSQL Server you will see the image below:

 ===========================================
= _ __ _____ ____ _ =
= ||/ // ____|/ __ \|| =
= |' /| (___ | | | | | = = | < \___ \| | | | | = = | . \ ____) | |__| | |____ = = |_|\_\_____/ \___\_\______| = = = = Streaming SQL Engine for Apache Kafka® = ===========================================Copyright 2017-2019 Confluent Inc.CLI v5.4.1, Server v5.4.1 located at http://localhost:8088Having trouble? Type 'help' (case-insensitive) for a rundown of how things work!ksql>

Some Commands

Show all topics

ksql> SHOW TOPICS;
Kafka Topic | Partitions | Partition Replicas
---------------------------------------------------------------------
com.mywebsite.streams.clickevents | 5 | 1
com.mywebsite.streams.pages | 1 | 1
---------------------------------------------------------------------

Show all streams

ksql> SHOW STREAMS;
Stream Name | Kafka Topic | Format
--------------------------------------------------------
CLICKEVENTS | com.mywebsite.streams.clickevents | JSON
--------------------------------------------------------

Creating a Stream

If you need run it in the background mode.

CREATE STREAM clickevents
(email VARCHAR,
timestampVARCHAR,
uri VARCHAR,
numberINTEGER)
WITH (KAFKA_TOPIC='com.mywebsite.streams.clickevents',
VALUE_FORMAT='JSON');

Creating a Table

CREATETABLEpages
(uri VARCHAR,
description VARCHAR,
created VARCHAR)
WITH (KAFKA_TOPIC='com.mywebsite.streams.pages',
VALUE_FORMAT='JSON',
KEY='uri');

Creating a Table from a Query

CREATETABLEa_pagesASSELECT*FROM pages WHERE uri LIKE'http://www.a%';

Querying a Table or Stream

SELECT*FROM clickevents EMIT CHANGES;

Describing a Table and Stream

ksql> DESCRIBE PAGES;
Name : PAGES
Field | Type
-----------------------------------------
ROWTIME | BIGINT (system)
ROWKEY | VARCHAR(STRING) (system)
URI | VARCHAR(STRING)
DESCRIPTION | VARCHAR(STRING)
CREATED | VARCHAR(STRING)
-----------------------------------------
For runtime statistics and query details run: DESCRIBE EXTENDED <Stream,Table>;

Managing Offsets

Like all Kafka Consumers, KSQL by default begins consumption at the latest offset. This can be a problem for some scenarios. In the following example we're going to create a pages table -- but -- we want all the data available to us in this table. In other words, we want KSQL to start from the earliest offset. To do this, we will use the SET command to set the configuration variabl auto.offset.reset for our session -- and before we run any commands.

SET 'auto.offset.reset' = 'earliest';

Also note that this can be set at the KSQL server level, if you'd like. Once you're done querying or creating tables or streams with this value, you can set it back to its original setting by simply running:

UNSET 'auto.offset.reset';

Scalar Functions

KSQL Provides a number of Scalar functions for us to make use of.

Lets write a function that takes advantage of some of these features:

SELECT UCASE(SUBSTRING(uri, 12))
FROM clickevents
WHERE number > 100
AND uri LIKE 'http://www.k%' EMIT CHANGES;

Notice that as soon as you hit CTRL+C your query ends

Deleting a Table

As with Streams, we must first find the running underlying query, and then drop the table. First, find your query:

ksql> SHOW QUERIES;
Query ID | Kafka Topic | Query String
----------------------------------------------------------------------------------------------
CTAS_A_PAGES_1 | A_PAGES | CREATE TABLE a_pages AS
SELECT * FROM pages WHERE uri LIKE 'http://www.a%';
----------------------------------------------------------------------------------------------
For detailed information on a Query run: EXPLAIN <Query ID>;

Find your query, which in this case is CTAS_A_PAGES_1 and then, finally, TERMINATE the query and DROP the table:

TERMINATE QUERY CTAS_A_PAGES_1;
DROP TABLE A_PAGES;

Windowing

Hopping and Tumbling Windows

In this demonstration we'll see how to create Tables with windowing enabled.

Tumbling Windows

Let's create a tumbling clickevents table, where the window size is 30 seconds.

CREATE STREAM clickevents_tumbling ASSELECT*FROM clickevents
WINDOW TUMBLING (SIZE 30 SECONDS);

Hopping Windows

Now we can create a Table with a hopping window of 30 seconds with 5 second increments.

CREATETABLEclickevents_hoppingASSELECT uri FROM clickevents
WINDOW HOPPING (SIZE 30 SECONDS, ADVANCE BY 5 SECONDS)
WHERE uri LIKE'http://www.b%'GROUP BY uri;

The above window is 30 seconds long and advances by 5 second. If you query the table you will see the associated window times!

Session Windows

Finally, lets see how session windows work. We're going to define the session as 5 minutes in order to group many events to the same window

CREATETABLEclickevents_sessionASSELECT uri FROM clickevents
WINDOW SESSION (5 MINUTES)
WHERE uri LIKE'http://www.b%'GROUP BY uri;

Kafka CLI Basic Commands

Creating a topic:

docker-compose exec kafka kafka-topics --create --topic <topic-name> --bootstrap-server localhost:9092

Writing a topic:

docker-compose exec kafka kafka-console-producer --topic <topic-name> --bootstrap-server localhost:9092

You must type on console the press key enter.

Reading a topic:

docker-compose exec kafka kafka-console-consumer.sh --topic <topic-name> --from-beginning --bootstrap-server localhost:9092

Contributing

Pull requests are welcome. For major changes, please open an issue first to discuss what you would like to change.

Please make sure to update tests as appropriate.

License

MIT

About

This project show how to use KSQL (Streaming SQL Engine for Apache Kafka) to stream processing.

Topics

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Latest commit

History

33 Commits

Folders and files

NameName
Last commit message
Last commit date

Repository files navigation

Stream Processing using KSQL

This project show how to use KSQL (Streaming SQL Engine for Apache Kafka) to stream processing.

Enviroment

For you to use this repository you will need the following softwares:

However, only Docker and Docker Compose need is installed in your machine. All Kafka ecosystem will be embedded via docker images.

Steps

  1. Install Python and Pip
  2. Install Docker and Docker Compose
  3. Load Images
  4. Create Topics
  5. Start Simulator

1 - Install Python and Pip

The installation process of the Python and Pip is very easy. So this tutorial dont't will cover this steps. I recommend you look for more information in www.python.org and pip.pypa.io.

After you install Python and Pip run the command below to install all dependencies need to execute click_simulator.py application. This code is responsible to simulate the click events into an web page. It will generate unbounded click events, sending a flow continuous messages to a Kafka topic.

pip install -r requirements.txt

2 - Install Docker and Docker Compose

This tutorial does not demonstrate the installation process for Docker and Docker Compose. I strongly recommend you to visit the Docker installation link for more informations. Please click here.

3 - Loading Images

docker-compose up

or

docker-compose up -d

The last command allow you to run docker-compose in the background.

4 - Create Topics

docker-compose exec kafka kafka-topics --create --topic com.mywebsite.streams.pages --bootstrap-server localhost:9092
docker-compose exec kafka kafka-topics --create --topic com.mywebsite.streams.clickevents --bootstrap-server localhost:9092

5 - Start Simulator

python click_simulator.py

If you executed all steps correctly. You will see an image similar that below.

Starting application
Message: {"email": "anoble@yahoo.com", "timestamp": "1986-03-10T16:38:40", "uri": "https://mitchell.info/login.php", "number": 358}
Message: {"email": "leonardpatrick@mason-clark.info", "timestamp": "1971-04-25T10:09:26", "uri": "https://www.bailey.com/search/about/", "number": 431}
Message: {"email": "morriskatie@villarreal-villa.biz", "timestamp": "1996-11-22T00:12:20", "uri": "http://www.woodard.info/terms.php", "number": 838}
Message: {"email": "kenneth79@rogers.info", "timestamp": "2005-10-24T22:16:59", "uri": "http://www.king.com/wp-content/blog/blog/index/", "number": 793}
Message: {"email": "wbailey@wu-martinez.net", "timestamp": "1995-06-20T12:44:44", "uri": "https://www.smith-neal.com/categories/login/", "number": 509}
Message: {"email": "tkennedy@hall-wolfe.org", "timestamp": "2009-01-27T14:04:20", "uri": "https://www.marshall-holmes.info/", "number": 336}
Message: {"email": "steven15@yahoo.com", "timestamp": "2019-12-13T16:09:11", "uri": "https://www.sims.net/main.html", "number": 263}
Message: {"email": "hobbsmario@hotmail.com", "timestamp": "1990-08-16T05:09:04", "uri": "http://www.smith.com/search/tags/explore/about.jsp", "number": 61}
...

Connecting to KSQL Server

docker-compose exec ksql ksql http://localhost:8088

After you connect to KSQL Server you will see the image below:

 ===========================================
= _ __ _____ ____ _ =
= ||/ // ____|/ __ \|| =
= |' /| (___ | | | | | = = | < \___ \| | | | | = = | . \ ____) | |__| | |____ = = |_|\_\_____/ \___\_\______| = = = = Streaming SQL Engine for Apache Kafka® = ===========================================Copyright 2017-2019 Confluent Inc.CLI v5.4.1, Server v5.4.1 located at http://localhost:8088Having trouble? Type 'help' (case-insensitive) for a rundown of how things work!ksql>

Some Commands

Show all topics

ksql> SHOW TOPICS;
Kafka Topic | Partitions | Partition Replicas
---------------------------------------------------------------------
com.mywebsite.streams.clickevents | 5 | 1
com.mywebsite.streams.pages | 1 | 1
---------------------------------------------------------------------

Show all streams

ksql> SHOW STREAMS;
Stream Name | Kafka Topic | Format
--------------------------------------------------------
CLICKEVENTS | com.mywebsite.streams.clickevents | JSON
--------------------------------------------------------

Creating a Stream

If you need run it in the background mode.

CREATE STREAM clickevents
(email VARCHAR,
timestampVARCHAR,
uri VARCHAR,
numberINTEGER)
WITH (KAFKA_TOPIC='com.mywebsite.streams.clickevents',
VALUE_FORMAT='JSON');

Creating a Table

CREATETABLEpages
(uri VARCHAR,
description VARCHAR,
created VARCHAR)
WITH (KAFKA_TOPIC='com.mywebsite.streams.pages',
VALUE_FORMAT='JSON',
KEY='uri');

Creating a Table from a Query

CREATETABLEa_pagesASSELECT*FROM pages WHERE uri LIKE'http://www.a%';

Querying a Table or Stream

SELECT*FROM clickevents EMIT CHANGES;

Describing a Table and Stream

ksql> DESCRIBE PAGES;
Name : PAGES
Field | Type
-----------------------------------------
ROWTIME | BIGINT (system)
ROWKEY | VARCHAR(STRING) (system)
URI | VARCHAR(STRING)
DESCRIPTION | VARCHAR(STRING)
CREATED | VARCHAR(STRING)
-----------------------------------------
For runtime statistics and query details run: DESCRIBE EXTENDED <Stream,Table>;

Managing Offsets

Like all Kafka Consumers, KSQL by default begins consumption at the latest offset. This can be a problem for some scenarios. In the following example we're going to create a pages table -- but -- we want all the data available to us in this table. In other words, we want KSQL to start from the earliest offset. To do this, we will use the SET command to set the configuration variabl auto.offset.reset for our session -- and before we run any commands.

SET 'auto.offset.reset' = 'earliest';

Also note that this can be set at the KSQL server level, if you'd like. Once you're done querying or creating tables or streams with this value, you can set it back to its original setting by simply running:

UNSET 'auto.offset.reset';

Scalar Functions

KSQL Provides a number of Scalar functions for us to make use of.

Lets write a function that takes advantage of some of these features:

SELECT UCASE(SUBSTRING(uri, 12))
FROM clickevents
WHERE number > 100
AND uri LIKE 'http://www.k%' EMIT CHANGES;

Notice that as soon as you hit CTRL+C your query ends

Deleting a Table

As with Streams, we must first find the running underlying query, and then drop the table. First, find your query:

ksql> SHOW QUERIES;
Query ID | Kafka Topic | Query String
----------------------------------------------------------------------------------------------
CTAS_A_PAGES_1 | A_PAGES | CREATE TABLE a_pages AS
SELECT * FROM pages WHERE uri LIKE 'http://www.a%';
----------------------------------------------------------------------------------------------
For detailed information on a Query run: EXPLAIN <Query ID>;

Find your query, which in this case is CTAS_A_PAGES_1 and then, finally, TERMINATE the query and DROP the table:

TERMINATE QUERY CTAS_A_PAGES_1;
DROP TABLE A_PAGES;

Windowing

Hopping and Tumbling Windows

In this demonstration we'll see how to create Tables with windowing enabled.

Tumbling Windows

Let's create a tumbling clickevents table, where the window size is 30 seconds.

CREATE STREAM clickevents_tumbling ASSELECT*FROM clickevents
WINDOW TUMBLING (SIZE 30 SECONDS);

Hopping Windows

Now we can create a Table with a hopping window of 30 seconds with 5 second increments.

CREATETABLEclickevents_hoppingASSELECT uri FROM clickevents
WINDOW HOPPING (SIZE 30 SECONDS, ADVANCE BY 5 SECONDS)
WHERE uri LIKE'http://www.b%'GROUP BY uri;

The above window is 30 seconds long and advances by 5 second. If you query the table you will see the associated window times!

Session Windows

Finally, lets see how session windows work. We're going to define the session as 5 minutes in order to group many events to the same window

CREATETABLEclickevents_sessionASSELECT uri FROM clickevents
WINDOW SESSION (5 MINUTES)
WHERE uri LIKE'http://www.b%'GROUP BY uri;

Kafka CLI Basic Commands

Creating a topic:

docker-compose exec kafka kafka-topics --create --topic <topic-name> --bootstrap-server localhost:9092

Writing a topic:

docker-compose exec kafka kafka-console-producer --topic <topic-name> --bootstrap-server localhost:9092

You must type on console the press key enter.

Reading a topic:

docker-compose exec kafka kafka-console-consumer.sh --topic <topic-name> --from-beginning --bootstrap-server localhost:9092

Contributing

Pull requests are welcome. For major changes, please open an issue first to discuss what you would like to change.

Please make sure to update tests as appropriate.

License

MIT

About

This project show how to use KSQL (Streaming SQL Engine for Apache Kafka) to stream processing.

Topics

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Latest commit

History

33 Commits

Folders and files

NameName
Last commit message
Last commit date

Repository files navigation

Stream Processing using KSQL

This project show how to use KSQL (Streaming SQL Engine for Apache Kafka) to stream processing.

Enviroment

For you to use this repository you will need the following softwares:

However, only Docker and Docker Compose need is installed in your machine. All Kafka ecosystem will be embedded via docker images.

Steps

  1. Install Python and Pip
  2. Install Docker and Docker Compose
  3. Load Images
  4. Create Topics
  5. Start Simulator

1 - Install Python and Pip

The installation process of the Python and Pip is very easy. So this tutorial dont't will cover this steps. I recommend you look for more information in www.python.org and pip.pypa.io.

After you install Python and Pip run the command below to install all dependencies need to execute click_simulator.py application. This code is responsible to simulate the click events into an web page. It will generate unbounded click events, sending a flow continuous messages to a Kafka topic.

pip install -r requirements.txt

2 - Install Docker and Docker Compose

This tutorial does not demonstrate the installation process for Docker and Docker Compose. I strongly recommend you to visit the Docker installation link for more informations. Please click here.

3 - Loading Images

docker-compose up

or

docker-compose up -d

The last command allow you to run docker-compose in the background.

4 - Create Topics

docker-compose exec kafka kafka-topics --create --topic com.mywebsite.streams.pages --bootstrap-server localhost:9092
docker-compose exec kafka kafka-topics --create --topic com.mywebsite.streams.clickevents --bootstrap-server localhost:9092

5 - Start Simulator

python click_simulator.py

If you executed all steps correctly. You will see an image similar that below.

Starting application
Message: {"email": "anoble@yahoo.com", "timestamp": "1986-03-10T16:38:40", "uri": "https://mitchell.info/login.php", "number": 358}
Message: {"email": "leonardpatrick@mason-clark.info", "timestamp": "1971-04-25T10:09:26", "uri": "https://www.bailey.com/search/about/", "number": 431}
Message: {"email": "morriskatie@villarreal-villa.biz", "timestamp": "1996-11-22T00:12:20", "uri": "http://www.woodard.info/terms.php", "number": 838}
Message: {"email": "kenneth79@rogers.info", "timestamp": "2005-10-24T22:16:59", "uri": "http://www.king.com/wp-content/blog/blog/index/", "number": 793}
Message: {"email": "wbailey@wu-martinez.net", "timestamp": "1995-06-20T12:44:44", "uri": "https://www.smith-neal.com/categories/login/", "number": 509}
Message: {"email": "tkennedy@hall-wolfe.org", "timestamp": "2009-01-27T14:04:20", "uri": "https://www.marshall-holmes.info/", "number": 336}
Message: {"email": "steven15@yahoo.com", "timestamp": "2019-12-13T16:09:11", "uri": "https://www.sims.net/main.html", "number": 263}
Message: {"email": "hobbsmario@hotmail.com", "timestamp": "1990-08-16T05:09:04", "uri": "http://www.smith.com/search/tags/explore/about.jsp", "number": 61}
...

Connecting to KSQL Server

docker-compose exec ksql ksql http://localhost:8088

After you connect to KSQL Server you will see the image below:

 ===========================================
= _ __ _____ ____ _ =
= ||/ // ____|/ __ \|| =
= |' /| (___ | | | | | = = | < \___ \| | | | | = = | . \ ____) | |__| | |____ = = |_|\_\_____/ \___\_\______| = = = = Streaming SQL Engine for Apache Kafka® = ===========================================Copyright 2017-2019 Confluent Inc.CLI v5.4.1, Server v5.4.1 located at http://localhost:8088Having trouble? Type 'help' (case-insensitive) for a rundown of how things work!ksql>

Some Commands

Show all topics

ksql> SHOW TOPICS;
Kafka Topic | Partitions | Partition Replicas
---------------------------------------------------------------------
com.mywebsite.streams.clickevents | 5 | 1
com.mywebsite.streams.pages | 1 | 1
---------------------------------------------------------------------

Show all streams

ksql> SHOW STREAMS;
Stream Name | Kafka Topic | Format
--------------------------------------------------------
CLICKEVENTS | com.mywebsite.streams.clickevents | JSON
--------------------------------------------------------

Creating a Stream

If you need run it in the background mode.

CREATE STREAM clickevents
(email VARCHAR,
timestampVARCHAR,
uri VARCHAR,
numberINTEGER)
WITH (KAFKA_TOPIC='com.mywebsite.streams.clickevents',
VALUE_FORMAT='JSON');

Creating a Table

CREATETABLEpages
(uri VARCHAR,
description VARCHAR,
created VARCHAR)
WITH (KAFKA_TOPIC='com.mywebsite.streams.pages',
VALUE_FORMAT='JSON',
KEY='uri');

Creating a Table from a Query

CREATETABLEa_pagesASSELECT*FROM pages WHERE uri LIKE'http://www.a%';

Querying a Table or Stream

SELECT*FROM clickevents EMIT CHANGES;

Describing a Table and Stream

ksql> DESCRIBE PAGES;
Name : PAGES
Field | Type
-----------------------------------------
ROWTIME | BIGINT (system)
ROWKEY | VARCHAR(STRING) (system)
URI | VARCHAR(STRING)
DESCRIPTION | VARCHAR(STRING)
CREATED | VARCHAR(STRING)
-----------------------------------------
For runtime statistics and query details run: DESCRIBE EXTENDED <Stream,Table>;

Managing Offsets

Like all Kafka Consumers, KSQL by default begins consumption at the latest offset. This can be a problem for some scenarios. In the following example we're going to create a pages table -- but -- we want all the data available to us in this table. In other words, we want KSQL to start from the earliest offset. To do this, we will use the SET command to set the configuration variabl auto.offset.reset for our session -- and before we run any commands.

SET 'auto.offset.reset' = 'earliest';

Also note that this can be set at the KSQL server level, if you'd like. Once you're done querying or creating tables or streams with this value, you can set it back to its original setting by simply running:

UNSET 'auto.offset.reset';

Scalar Functions

KSQL Provides a number of Scalar functions for us to make use of.

Lets write a function that takes advantage of some of these features:

SELECT UCASE(SUBSTRING(uri, 12))
FROM clickevents
WHERE number > 100
AND uri LIKE 'http://www.k%' EMIT CHANGES;

Notice that as soon as you hit CTRL+C your query ends

Deleting a Table

As with Streams, we must first find the running underlying query, and then drop the table. First, find your query:

ksql> SHOW QUERIES;
Query ID | Kafka Topic | Query String
----------------------------------------------------------------------------------------------
CTAS_A_PAGES_1 | A_PAGES | CREATE TABLE a_pages AS
SELECT * FROM pages WHERE uri LIKE 'http://www.a%';
----------------------------------------------------------------------------------------------
For detailed information on a Query run: EXPLAIN <Query ID>;

Find your query, which in this case is CTAS_A_PAGES_1 and then, finally, TERMINATE the query and DROP the table:

TERMINATE QUERY CTAS_A_PAGES_1;
DROP TABLE A_PAGES;

Windowing

Hopping and Tumbling Windows

In this demonstration we'll see how to create Tables with windowing enabled.

Tumbling Windows

Let's create a tumbling clickevents table, where the window size is 30 seconds.

CREATE STREAM clickevents_tumbling ASSELECT*FROM clickevents
WINDOW TUMBLING (SIZE 30 SECONDS);

Hopping Windows

Now we can create a Table with a hopping window of 30 seconds with 5 second increments.

CREATETABLEclickevents_hoppingASSELECT uri FROM clickevents
WINDOW HOPPING (SIZE 30 SECONDS, ADVANCE BY 5 SECONDS)
WHERE uri LIKE'http://www.b%'GROUP BY uri;

The above window is 30 seconds long and advances by 5 second. If you query the table you will see the associated window times!

Session Windows

Finally, lets see how session windows work. We're going to define the session as 5 minutes in order to group many events to the same window

CREATETABLEclickevents_sessionASSELECT uri FROM clickevents
WINDOW SESSION (5 MINUTES)
WHERE uri LIKE'http://www.b%'GROUP BY uri;

Kafka CLI Basic Commands

Creating a topic:

docker-compose exec kafka kafka-topics --create --topic <topic-name> --bootstrap-server localhost:9092

Writing a topic:

docker-compose exec kafka kafka-console-producer --topic <topic-name> --bootstrap-server localhost:9092

You must type on console the press key enter.

Reading a topic:

docker-compose exec kafka kafka-console-consumer.sh --topic <topic-name> --from-beginning --bootstrap-server localhost:9092

Contributing

Pull requests are welcome. For major changes, please open an issue first to discuss what you would like to change.

Please make sure to update tests as appropriate.

License

MIT

About

This project show how to use KSQL (Streaming SQL Engine for Apache Kafka) to stream processing.

Topics

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content

Latest commit

History

33 Commits

Folders and files

NameName
Last commit message
Last commit date

Repository files navigation

Stream Processing using KSQL

This project show how to use KSQL (Streaming SQL Engine for Apache Kafka) to stream processing.

Enviroment

For you to use this repository you will need the following softwares:

However, only Docker and Docker Compose need is installed in your machine. All Kafka ecosystem will be embedded via docker images.

Steps

  1. Install Python and Pip
  2. Install Docker and Docker Compose
  3. Load Images
  4. Create Topics
  5. Start Simulator

1 - Install Python and Pip

The installation process of the Python and Pip is very easy. So this tutorial dont't will cover this steps. I recommend you look for more information in www.python.org and pip.pypa.io.

After you install Python and Pip run the command below to install all dependencies need to execute click_simulator.py application. This code is responsible to simulate the click events into an web page. It will generate unbounded click events, sending a flow continuous messages to a Kafka topic.

pip install -r requirements.txt

2 - Install Docker and Docker Compose

This tutorial does not demonstrate the installation process for Docker and Docker Compose. I strongly recommend you to visit the Docker installation link for more informations. Please click here.

3 - Loading Images

docker-compose up

or

docker-compose up -d

The last command allow you to run docker-compose in the background.

4 - Create Topics

docker-compose exec kafka kafka-topics --create --topic com.mywebsite.streams.pages --bootstrap-server localhost:9092
docker-compose exec kafka kafka-topics --create --topic com.mywebsite.streams.clickevents --bootstrap-server localhost:9092

5 - Start Simulator

python click_simulator.py

If you executed all steps correctly. You will see an image similar that below.

Starting application
Message: {"email": "anoble@yahoo.com", "timestamp": "1986-03-10T16:38:40", "uri": "https://mitchell.info/login.php", "number": 358}
Message: {"email": "leonardpatrick@mason-clark.info", "timestamp": "1971-04-25T10:09:26", "uri": "https://www.bailey.com/search/about/", "number": 431}
Message: {"email": "morriskatie@villarreal-villa.biz", "timestamp": "1996-11-22T00:12:20", "uri": "http://www.woodard.info/terms.php", "number": 838}
Message: {"email": "kenneth79@rogers.info", "timestamp": "2005-10-24T22:16:59", "uri": "http://www.king.com/wp-content/blog/blog/index/", "number": 793}
Message: {"email": "wbailey@wu-martinez.net", "timestamp": "1995-06-20T12:44:44", "uri": "https://www.smith-neal.com/categories/login/", "number": 509}
Message: {"email": "tkennedy@hall-wolfe.org", "timestamp": "2009-01-27T14:04:20", "uri": "https://www.marshall-holmes.info/", "number": 336}
Message: {"email": "steven15@yahoo.com", "timestamp": "2019-12-13T16:09:11", "uri": "https://www.sims.net/main.html", "number": 263}
Message: {"email": "hobbsmario@hotmail.com", "timestamp": "1990-08-16T05:09:04", "uri": "http://www.smith.com/search/tags/explore/about.jsp", "number": 61}
...

Connecting to KSQL Server

docker-compose exec ksql ksql http://localhost:8088

After you connect to KSQL Server you will see the image below:

 ===========================================
= _ __ _____ ____ _ =
= ||/ // ____|/ __ \|| =
= |' /| (___ | | | | | = = | < \___ \| | | | | = = | . \ ____) | |__| | |____ = = |_|\_\_____/ \___\_\______| = = = = Streaming SQL Engine for Apache Kafka® = ===========================================Copyright 2017-2019 Confluent Inc.CLI v5.4.1, Server v5.4.1 located at http://localhost:8088Having trouble? Type 'help' (case-insensitive) for a rundown of how things work!ksql>

Some Commands

Show all topics

ksql> SHOW TOPICS;
Kafka Topic | Partitions | Partition Replicas
---------------------------------------------------------------------
com.mywebsite.streams.clickevents | 5 | 1
com.mywebsite.streams.pages | 1 | 1
---------------------------------------------------------------------

Show all streams

ksql> SHOW STREAMS;
Stream Name | Kafka Topic | Format
--------------------------------------------------------
CLICKEVENTS | com.mywebsite.streams.clickevents | JSON
--------------------------------------------------------

Creating a Stream

If you need run it in the background mode.

CREATE STREAM clickevents
(email VARCHAR,
timestampVARCHAR,
uri VARCHAR,
numberINTEGER)
WITH (KAFKA_TOPIC='com.mywebsite.streams.clickevents',
VALUE_FORMAT='JSON');

Creating a Table

CREATETABLEpages
(uri VARCHAR,
description VARCHAR,
created VARCHAR)
WITH (KAFKA_TOPIC='com.mywebsite.streams.pages',
VALUE_FORMAT='JSON',
KEY='uri');

Creating a Table from a Query

CREATETABLEa_pagesASSELECT*FROM pages WHERE uri LIKE'http://www.a%';

Querying a Table or Stream

SELECT*FROM clickevents EMIT CHANGES;

Describing a Table and Stream

ksql> DESCRIBE PAGES;
Name : PAGES
Field | Type
-----------------------------------------
ROWTIME | BIGINT (system)
ROWKEY | VARCHAR(STRING) (system)
URI | VARCHAR(STRING)
DESCRIPTION | VARCHAR(STRING)
CREATED | VARCHAR(STRING)
-----------------------------------------
For runtime statistics and query details run: DESCRIBE EXTENDED <Stream,Table>;

Managing Offsets

Like all Kafka Consumers, KSQL by default begins consumption at the latest offset. This can be a problem for some scenarios. In the following example we're going to create a pages table -- but -- we want all the data available to us in this table. In other words, we want KSQL to start from the earliest offset. To do this, we will use the SET command to set the configuration variabl auto.offset.reset for our session -- and before we run any commands.

SET 'auto.offset.reset' = 'earliest';

Also note that this can be set at the KSQL server level, if you'd like. Once you're done querying or creating tables or streams with this value, you can set it back to its original setting by simply running:

UNSET 'auto.offset.reset';

Scalar Functions

KSQL Provides a number of Scalar functions for us to make use of.

Lets write a function that takes advantage of some of these features:

SELECT UCASE(SUBSTRING(uri, 12))
FROM clickevents
WHERE number > 100
AND uri LIKE 'http://www.k%' EMIT CHANGES;

Notice that as soon as you hit CTRL+C your query ends

Deleting a Table

As with Streams, we must first find the running underlying query, and then drop the table. First, find your query:

ksql> SHOW QUERIES;
Query ID | Kafka Topic | Query String
----------------------------------------------------------------------------------------------
CTAS_A_PAGES_1 | A_PAGES | CREATE TABLE a_pages AS
SELECT * FROM pages WHERE uri LIKE 'http://www.a%';
----------------------------------------------------------------------------------------------
For detailed information on a Query run: EXPLAIN <Query ID>;

Find your query, which in this case is CTAS_A_PAGES_1 and then, finally, TERMINATE the query and DROP the table:

TERMINATE QUERY CTAS_A_PAGES_1;
DROP TABLE A_PAGES;

Windowing

Hopping and Tumbling Windows

In this demonstration we'll see how to create Tables with windowing enabled.

Tumbling Windows

Let's create a tumbling clickevents table, where the window size is 30 seconds.

CREATE STREAM clickevents_tumbling ASSELECT*FROM clickevents
WINDOW TUMBLING (SIZE 30 SECONDS);

Hopping Windows

Now we can create a Table with a hopping window of 30 seconds with 5 second increments.

CREATETABLEclickevents_hoppingASSELECT uri FROM clickevents
WINDOW HOPPING (SIZE 30 SECONDS, ADVANCE BY 5 SECONDS)
WHERE uri LIKE'http://www.b%'GROUP BY uri;

The above window is 30 seconds long and advances by 5 second. If you query the table you will see the associated window times!

Session Windows

Finally, lets see how session windows work. We're going to define the session as 5 minutes in order to group many events to the same window

CREATETABLEclickevents_sessionASSELECT uri FROM clickevents
WINDOW SESSION (5 MINUTES)
WHERE uri LIKE'http://www.b%'GROUP BY uri;

Kafka CLI Basic Commands

Creating a topic:

docker-compose exec kafka kafka-topics --create --topic <topic-name> --bootstrap-server localhost:9092

Writing a topic:

docker-compose exec kafka kafka-console-producer --topic <topic-name> --bootstrap-server localhost:9092

You must type on console the press key enter.

Reading a topic:

docker-compose exec kafka kafka-console-consumer.sh --topic <topic-name> --from-beginning --bootstrap-server localhost:9092

Contributing

Pull requests are welcome. For major changes, please open an issue first to discuss what you would like to change.

Please make sure to update tests as appropriate.

License

MIT

About

This project show how to use KSQL (Streaming SQL Engine for Apache Kafka) to stream processing.

Topics

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages