Have you checked out my YouTube Channel yet? See all these posts demonstrated from end-to-end

Oracle 26ai and Kafka: Ingest Kafka Events Directly into Oracle with DBMS_KAFKA

Your vaccine shipment just crossed 8°C. The reading is already sitting in Kafka. Can the Oracle database see that breach without you writing a Kafka consumer application?

Oracle 26ai and Kafka: Ingest Kafka Events Directly into Oracle with DBMS_KAFKA

In this post we connect Apache Kafka to Oracle AI Database 26ai Free (both in Docker), pull sensor events into a table with the built-in DBMS_KAFKA package, and let the database flag temperature breaches on its own. Every command below was run, in this order, from a clean state.

Temperature sensor
      |
      v
 Apache Kafka  (topic: coldchain.readings)
      |
      |  DBMS_KAFKA  (inside Oracle)
      v
 COLDCHAIN_READINGS   (raw Kafka JSON)
      |
      v
 JSON_TABLE + shipment rules   -->   OK / TOO_HOT / TOO_COLD

A fair warning up front: this feature (Oracle SQL Access to Kafka) is not new in 26ai. It first appeared in 23ai, and 26ai Free carries it forward. What is uncommon is running everything locally, with your own Kafka broker in Docker instead of a cloud streaming service.

How to read this post: three places you run things

Almost every mistake in a setup like this is running the right command in the wrong place. So every code block below is labelled with exactly where to run it:

  • 🟥 SYS — the SQLcl connection named Local Sys DBA (administration: grants, directories, registering Kafka).
  • 🟦 AIUSER — the SQLcl connection named Docker AI (all the application objects).
  • ⬛ TERMINAL — a normal terminal on your machine (Docker and Kafka commands).

Both SQL connections point at localhost:1521/FREEPDB1. For multi-statement blocks use Run Script (not Run Statement) in your SQL editor.

The whole sequence at a glance

StepWhatWhere
0Pre-flight checksTERMINAL, AIUSER, SYS
1Confirm Kafka support exists in the databaseAIUSER
2Prove Oracle can reach KafkaTERMINAL
3Create the Kafka topicTERMINAL
4Register the Kafka cluster in OracleTERMINAL, then SYS
5Create the load applicationAIUSER
6The ORA-62828 gotchaAIUSER
7Shipment rules and the business viewAIUSER
8First messages and manual loadTERMINAL, then AIUSER
9Automate the load with a scheduler jobAIUSER
10The payoff: Kafka message to TOO_HOTAIUSER, TERMINAL, AIUSER
11What just happened?AIUSER
12Prove the rules (OK, TOO_HOT, TOO_COLD)TERMINAL, then AIUSER
 Troubleshooting, reset, build from scratchvarious

Step 0: Pre-flight checks

This post assumes two containers are already running: the Oracle 26ai Free container (oracle26ai-apex261) and a Kafka container (kafka), both attached to a Docker network called kafka-net. If you don't have them yet, jump to Build from scratch at the end.

⬛ Run in: TERMINAL — PowerShell / Command Prompt on your machine
Both containers should be Up, and Oracle should say (healthy).

docker ps --format "table {{.Names}}\t{{.Status}}\t{{.Ports}}"

⬛ Run in: TERMINAL — PowerShell / Command Prompt on your machine
Both container names should appear in the output.

docker network inspect kafka-net --format "{{range .Containers}}{{.Name}} {{end}}"

Now make sure no leftovers from earlier experiments exist in the database. Otherwise you will get “object already exists” errors halfway through.

🟦 Run as: AIUSER — SQLcl connection Docker AI
Expect 0 for every row.

select 'OBJECTS' what, count(*) n from user_objects
 where object_name like 'COLDCHAIN%' or object_name like 'ORA$DKV%'
union all
select 'JOBS', count(*) from user_scheduler_jobs where job_name like 'COLDCHAIN%'
union all
select 'KAFKA_APPS', count(*) from user_kafka_applications;

🟥 Run as: SYS — SQLcl connection Local Sys DBA
Expect 0 for every row.

select 'CLUSTERS' what, count(*) n from dbms_kafka_clusters
union all
select 'DIRECTORIES', count(*) from dba_directories where directory_name like 'OSAK_COLD%'
union all
select 'AIUSER_OSAK_ROLE', count(*) from dba_role_privs
 where grantee = 'AIUSER' and granted_role = 'OSAK_ADMIN_ROLE';

Not zero? Use the Reset section at the end, then come back here.

Step 1: Does this database already speak Kafka?

Before configuring anything, check that Kafka support exists. This is not a library you install: Oracle SQL Access to Kafka (OSAK) is part of the database, and DBMS_KAFKA is its PL/SQL API.

🟦 Run as: AIUSER — SQLcl connection Docker AI
Connected as AIUSER.

select owner, object_type, object_name
from   all_objects
where  object_name = 'DBMS_KAFKA'
and    object_type = 'PACKAGE';

🟦 Run as: AIUSER — SQLcl connection Docker AI
The API surface: load apps, streaming apps, seekable apps, offsets.

select distinct procedure_name
from   all_procedures
where  object_name = 'DBMS_KAFKA'
and    procedure_name is not null
order  by 1;

You will see CREATE_LOAD_APP, EXECUTE_LOAD_APP, CREATE_STREAMING_APP and friends. Notice what is not there: nothing that publishes to Kafka. DBMS_KAFKA is the Kafka-to-Oracle direction only. Keep that in mind for the end of the post.

Step 2: Can Oracle actually reach Kafka?

Oracle runs inside a container, so localhost means the Oracle container itself, not your PC. From inside the Docker network, Kafka is at kafka:9092. From your Windows host it is localhost:29092. Two addresses, one broker. Test the one Oracle will use:

⬛ Run in: TERMINAL — PowerShell / Command Prompt on your machine
Expect REACHABLE. If you see UNREACHABLE, stop and fix networking first (see Troubleshooting).

docker exec oracle26ai-apex261 bash -c "timeout 5 bash -c '</dev/tcp/kafka/9092' && echo REACHABLE || echo UNREACHABLE"

Step 3: Create the topic

Three partitions make it look like a realistic event stream. --if-not-exists makes the command safe to run again.

⬛ Run in: TERMINAL — PowerShell / Command Prompt on your machine
Expect “Created topic coldchain.readings” (the metric-name warning is harmless).

docker exec kafka /opt/kafka/bin/kafka-topics.sh --bootstrap-server kafka:9092 --create --if-not-exists --topic coldchain.readings --partitions 3 --replication-factor 1

Step 4: Introduce Kafka to Oracle

Oracle needs to be told about the Kafka cluster. That takes an OS folder, two directory objects, an admin role, and one registration call. Note the switch of users: this step is mostly SYS.

4a. Create the config folder inside the Oracle container

For an unsecured (PLAINTEXT) broker, the required osakafka.properties file can be empty.

⬛ Run in: TERMINAL — PowerShell / Command Prompt on your machine
Expect to see an empty osakafka.properties file listed.

docker exec oracle26ai-apex261 bash -c "mkdir -p /opt/oracle/osak/coldchain/config && touch /opt/oracle/osak/coldchain/config/osakafka.properties && ls -l /opt/oracle/osak/coldchain/config"

4b. Grants and directories

Switch to SYS. We give AIUSER the OSAK admin role and permission to create scheduler jobs (needed later in Step 9), and create the two directory objects OSAK requires.

🟥 Run as: SYS — SQLcl connection Local Sys DBA
Run as a script. Every statement should succeed.

grant OSAK_ADMIN_ROLE to aiuser;
grant create job to aiuser;

create or replace directory OSAK_COLDCHAIN_CONFIG as '/opt/oracle/osak/coldchain/config';
create or replace directory OSAK_COLDCHAIN_ACCESS as '/opt/oracle/osak/coldchain/config';

grant read on directory OSAK_COLDCHAIN_CONFIG to aiuser;
grant read on directory OSAK_COLDCHAIN_ACCESS to aiuser;

4c. Register the cluster

This is the moment Oracle learns about your Kafka. REGISTER_CLUSTER is a function, and it tests the connection itself: a return value of 0 means success.

🟥 Run as: SYS — SQLcl connection Local Sys DBA
Expect RC = 0.

select DBMS_KAFKA_ADM.REGISTER_CLUSTER(
         cluster_name        => 'COLDCHAIN',
         bootstrap_servers   => 'kafka:9092',
         kafka_provider      => 'APACHE',
         cluster_access_dir  => 'OSAK_COLDCHAIN_ACCESS',
         credential_name     => NULL,
         cluster_config_dir  => 'OSAK_COLDCHAIN_CONFIG',
         cluster_description => 'Cold chain demo Kafka (Docker)',
         options             => NULL) as rc
from dual;

4d. Checkpoint

🟥 Run as: SYS — SQLcl connection Local Sys DBA
Expect COLDCHAIN | 0 | kafka:9092 | APACHE. State 0 means connected.

select cluster_name, state, bootstrap_servers, kafka_provider
from   dbms_kafka_clusters;

Switch back to AIUSER now. Everything from here on, until the troubleshooting and reset sections, runs as AIUSER (plus the Terminal for producing messages).

Step 5: Create the load application

An OSAK load application reads a topic incrementally, from where it last stopped up to the current end of the topic, and lands the records in an ordinary Oracle table. The option fmt=JSON says our messages are JSON. One quirk: application names allow only letters, digits and $. No underscores.

🟦 Run as: AIUSER — SQLcl connection Docker AI
Expect “PL/SQL procedure successfully completed”.

begin
  DBMS_KAFKA.CREATE_LOAD_APP(
    cluster_name     => 'COLDCHAIN',
    application_name => 'CCLOAD',
    topic_name       => 'coldchain.readings',
    options          => '{"fmt":"JSON"}');
end;
/

🟦 Run as: AIUSER — SQLcl connection Docker AI
Expect one row: CCLOAD | COLDCHAIN | coldchain.readings | LOAD.

select application_name, cluster_name, topic_name, application_type
from   user_kafka_applications;

Rule of thumb: do not query the generated ORA$DKV_... view yourself. OSAK dedicates it to the application, and querying it directly can interfere with the offsets.

Step 6: The ORA-62828 gotcha (do this on purpose)

The natural thing is to create a target table with your own extra columns, such as a loaded_at timestamp. Try it:

🟦 Run as: AIUSER — SQLcl connection Docker AI
The table creates fine. The load will not.

create table coldchain_readings (
  kafka_partition       number,
  kafka_offset          number,
  kafka_epoch_timestamp number,
  value                 varchar2(4000),
  loaded_at             timestamp default systimestamp
);

set serveroutput on
declare
  n integer;
begin
  DBMS_KAFKA.EXECUTE_LOAD_APP('COLDCHAIN', 'CCLOAD', 'COLDCHAIN_READINGS', n);
  dbms_output.put_line('records_loaded=' || n);
end;
/

You get ORA-62828: the target table contains one or more column names that do not exist in the OSAK view. Every column in the target must exist in the view OSAK generates, which has just four: partition, offset, epoch timestamp and value. The lesson is architectural: keep the ingestion table as a faithful copy of what Kafka delivered, and put business meaning in a view on top.

🟦 Run as: AIUSER — SQLcl connection Docker AI
Drop the wrong table and create the right one.

drop table coldchain_readings purge;

create table coldchain_readings (
  kafka_partition       number         not null,
  kafka_offset          number         not null,
  kafka_epoch_timestamp number,
  value                 varchar2(4000) constraint cc_readings_json_ck check (value is json),
  constraint coldchain_readings_pk primary key (kafka_partition, kafka_offset)
);

Step 7: Kafka knows the temperature. Oracle knows what is acceptable.

A reading of 9°C is a disaster for vaccines and wonderful for frozen seafood. Kafka cannot tell the difference; Oracle can, if we give it the shipment rules.

🟦 Run as: AIUSER — SQLcl connection Docker AI
Reference data: the safe temperature range per shipment.

create table coldchain_shipments (
  shipment_id  varchar2(30) primary key,
  product      varchar2(100) not null,
  origin       varchar2(60),
  destination  varchar2(60),
  min_temp_c   number(5,2) not null,
  max_temp_c   number(5,2) not null
);

insert into coldchain_shipments values ('SHP-1001', 'COVID/Flu vaccines', 'Pune',      'Chennai',   2,   8);
insert into coldchain_shipments values ('SHP-1002', 'Frozen seafood',     'Kochi',     'Dubai',   -25, -15);
insert into coldchain_shipments values ('SHP-1003', 'Insulin pens',        'Hyderabad', 'Bengaluru', 2,   8);
commit;

select * from coldchain_shipments order by shipment_id;

Here is one raw Kafka message, then the view that turns that JSON into columns with JSON_TABLE and joins it to the rules:

{"sensor_id":"S-001","shipment_id":"SHP-1001","temp_c":4.2,"humidity_pct":61.5,"ts":"2026-10-03T06:30:00"}

🟦 Run as: AIUSER — SQLcl connection Docker AI
Expect “View COLDCHAIN_READINGS_V created”.

create or replace view coldchain_readings_v as
select r.kafka_partition,
       r.kafka_offset,
       j.sensor_id,
       j.shipment_id,
       j.temp_c,
       j.humidity_pct,
       to_timestamp(j.ts, 'YYYY-MM-DD"T"HH24:MI:SS') as reading_ts,
       s.product,
       s.min_temp_c,
       s.max_temp_c,
       case when j.temp_c > s.max_temp_c then 'TOO_HOT'
            when j.temp_c < s.min_temp_c then 'TOO_COLD'
            else 'OK'
       end as status
from   coldchain_readings r
cross apply json_table(r.value, '$' columns (
         sensor_id    varchar2(30) path '$.sensor_id',
         shipment_id  varchar2(30) path '$.shipment_id',
         temp_c       number       path '$.temp_c',
         humidity_pct number       path '$.humidity_pct',
         ts           varchar2(30) path '$.ts')) j
left join coldchain_shipments s on s.shipment_id = j.shipment_id;

Step 8: Send the first events

Time for data. The easiest quote-proof way on Windows is the interactive producer: start it, paste JSON lines at the > prompt (one line = one Kafka message), then press Ctrl+C to leave.

⬛ Run in: TERMINAL — PowerShell / Command Prompt on your machine
Starts the producer. It waits at a > prompt.

docker exec -it kafka /opt/kafka/bin/kafka-console-producer.sh --bootstrap-server kafka:9092 --topic coldchain.readings

⬛ Run in: TERMINAL — PowerShell / Command Prompt on your machine
Paste this at the > prompt, press Enter, then Ctrl+C. (If your terminal asks about multi-line paste, allow it.)

{"sensor_id":"S-001","shipment_id":"SHP-1001","temp_c":4.2,"humidity_pct":61.5,"ts":"2026-10-03T06:30:00"}

Prefer a one-liner? This worked in a Windows Command Prompt (cmd). PowerShell handles the escaped quotes differently, so stay with the interactive method there:

docker exec kafka sh -c "echo '{\"sensor_id\":\"S-001\",\"shipment_id\":\"SHP-1001\",\"temp_c\":4.2,\"humidity_pct\":61.5,\"ts\":\"2026-10-03T06:30:00\"}' | /opt/kafka/bin/kafka-console-producer.sh --bootstrap-server kafka:9092 --topic coldchain.readings"

Optional: prove the message really is in Kafka before Oracle touches it.

⬛ Run in: TERMINAL — PowerShell / Command Prompt on your machine
Prints the message and “Processed a total of 1 messages”.

docker exec kafka /opt/kafka/bin/kafka-console-consumer.sh --bootstrap-server kafka:9092 --topic coldchain.readings --from-beginning --max-messages 1 --timeout-ms 8000

Now pull it into Oracle manually, once. Each call to EXECUTE_LOAD_APP loads only the records it has not loaded before.

🟦 Run as: AIUSER — SQLcl connection Docker AI
Expect one row: SHP-1001 | 4.2 | OK.

set serveroutput on
declare
  n integer;
begin
  DBMS_KAFKA.EXECUTE_LOAD_APP('COLDCHAIN', 'CCLOAD', 'COLDCHAIN_READINGS', n);
  commit;
  dbms_output.put_line('records_loaded=' || n);
end;
/

select shipment_id, temp_c, status from coldchain_readings_v;

A realistic baseline: 12 more readings

Three shipments, with vaccines creeping above their limit and the seafood warming up. Start the producer again (same command as above), paste all 12 lines, press Enter, then Ctrl+C.

⬛ Run in: TERMINAL — PowerShell / Command Prompt on your machine
Producer command (same as before).

docker exec -it kafka /opt/kafka/bin/kafka-console-producer.sh --bootstrap-server kafka:9092 --topic coldchain.readings

⬛ Run in: TERMINAL — PowerShell / Command Prompt on your machine
Paste all 12 lines at the > prompt.

{"sensor_id":"S-001","shipment_id":"SHP-1001","temp_c":4.5,"humidity_pct":60.1,"ts":"2026-10-03T06:40:00"}
{"sensor_id":"S-001","shipment_id":"SHP-1001","temp_c":5.1,"humidity_pct":60.1,"ts":"2026-10-03T06:41:00"}
{"sensor_id":"S-001","shipment_id":"SHP-1001","temp_c":6.8,"humidity_pct":60.1,"ts":"2026-10-03T06:42:00"}
{"sensor_id":"S-001","shipment_id":"SHP-1001","temp_c":8.4,"humidity_pct":60.1,"ts":"2026-10-03T06:43:00"}
{"sensor_id":"S-001","shipment_id":"SHP-1001","temp_c":9.7,"humidity_pct":60.1,"ts":"2026-10-03T06:44:00"}
{"sensor_id":"S-002","shipment_id":"SHP-1002","temp_c":-21.5,"humidity_pct":48.0,"ts":"2026-10-03T06:45:00"}
{"sensor_id":"S-002","shipment_id":"SHP-1002","temp_c":-19.8,"humidity_pct":48.0,"ts":"2026-10-03T06:46:00"}
{"sensor_id":"S-002","shipment_id":"SHP-1002","temp_c":-16.2,"humidity_pct":48.0,"ts":"2026-10-03T06:47:00"}
{"sensor_id":"S-002","shipment_id":"SHP-1002","temp_c":-13.9,"humidity_pct":48.0,"ts":"2026-10-03T06:48:00"}
{"sensor_id":"S-003","shipment_id":"SHP-1003","temp_c":5.0,"humidity_pct":55.0,"ts":"2026-10-03T06:50:00"}
{"sensor_id":"S-003","shipment_id":"SHP-1003","temp_c":5.2,"humidity_pct":55.0,"ts":"2026-10-03T06:51:00"}
{"sensor_id":"S-003","shipment_id":"SHP-1003","temp_c":4.9,"humidity_pct":55.0,"ts":"2026-10-03T06:52:00"}

🟦 Run as: AIUSER — SQLcl connection Docker AI
Load them manually one last time (this is the last manual load).

declare
  n integer;
begin
  DBMS_KAFKA.EXECUTE_LOAD_APP('COLDCHAIN', 'CCLOAD', 'COLDCHAIN_READINGS', n);
  commit;
end;
/

select shipment_id, product, count(*) readings, min(temp_c) min_t, max(temp_c) max_t,
       count(case when status <> 'OK' then 1 end) breaches
from   coldchain_readings_v
group  by shipment_id, product
order  by shipment_id;

Expected: SHP-1001 has 6 readings and 2 breaches, SHP-1002 has 4 readings and 1 breach, SHP-1003 has 3 readings and none. That is 13 rows and 3 breaches in total.

Step 9: Stop pressing the button

A stream you must trigger by hand isn't a stream. We wrap the load in a small procedure, log each run that actually loaded something, and let DBMS_SCHEDULER call it every 5 seconds. (This is why Step 4b granted CREATE JOB.)

🟦 Run as: AIUSER — SQLcl connection Docker AI
Run as a script. Expect the job to show ENABLED = TRUE at the end.

create table coldchain_load_log (
  run_id          number generated always as identity primary key,
  run_ts          timestamp default systimestamp not null,
  records_loaded  number,
  status          varchar2(10),
  error_msg       varchar2(4000)
);

create or replace procedure coldchain_load_kafka as
  n   integer := 0;
  msg varchar2(4000);
begin
  dbms_kafka.execute_load_app('COLDCHAIN', 'CCLOAD', 'COLDCHAIN_READINGS', n);
  commit;
  if n > 0 then
    insert into coldchain_load_log(records_loaded, status) values (n, 'OK');
    commit;
  end if;
exception
  when others then
    msg := substr(sqlerrm, 1, 4000);
    rollback;
    insert into coldchain_load_log(records_loaded, status, error_msg) values (0, 'ERROR', msg);
    commit;
end;
/

begin
  dbms_scheduler.create_job(
    job_name        => 'COLDCHAIN_KAFKA_LOAD_JOB',
    job_type        => 'STORED_PROCEDURE',
    job_action      => 'COLDCHAIN_LOAD_KAFKA',
    repeat_interval => 'FREQ=SECONDLY;INTERVAL=5',
    enabled         => TRUE,
    comments        => 'Pulls new Kafka records from coldchain.readings into COLDCHAIN_READINGS');
end;
/

select job_name, enabled, state from user_scheduler_jobs where job_name = 'COLDCHAIN_KAFKA_LOAD_JOB';

Step 10: The payoff

Here is what viewers should remember: Kafka message, a few seconds, Oracle query, TOO_HOT. We will not call EXECUTE_LOAD_APP this time.

🟦 Run as: AIUSER — SQLcl connection Docker AI
Before: expect TOTAL_ROWS = 13 and BREACHES = 3.

select (select count(*) from coldchain_readings) total_rows,
       (select count(*) from coldchain_readings_v where status <> 'OK') breaches
from dual;

⬛ Run in: TERMINAL — PowerShell / Command Prompt on your machine
Start the producer.

docker exec -it kafka /opt/kafka/bin/kafka-console-producer.sh --bootstrap-server kafka:9092 --topic coldchain.readings

⬛ Run in: TERMINAL — PowerShell / Command Prompt on your machine
Three vaccine readings, all above the 8°C limit. Paste, Enter, Ctrl+C.

{"sensor_id":"S-001","shipment_id":"SHP-1001","temp_c":9.1,"humidity_pct":63.0,"ts":"2026-10-03T07:10:00"}
{"sensor_id":"S-001","shipment_id":"SHP-1001","temp_c":10.3,"humidity_pct":63.0,"ts":"2026-10-03T07:11:00"}
{"sensor_id":"S-001","shipment_id":"SHP-1001","temp_c":11.8,"humidity_pct":63.0,"ts":"2026-10-03T07:12:00"}

Now wait about 5 to 10 seconds (the job runs every 5 seconds), then run:

🟦 Run as: AIUSER — SQLcl connection Docker AI
Expect 11.8, 10.3 and 9.1 at the top, each TOO_HOT.

select reading_ts, shipment_id, product, temp_c, status
from   coldchain_readings_v
order  by reading_ts desc, kafka_offset desc
fetch  first 10 rows only;

🟦 Run as: AIUSER — SQLcl connection Docker AI
After: expect TOTAL_ROWS = 16 and BREACHES = 6.

select (select count(*) from coldchain_readings) total_rows,
       (select count(*) from coldchain_readings_v where status <> 'OK') breaches
from dual;

Nobody ran the load. Oracle noticed on its own.

Step 11: So what actually happened?

🟦 Run as: AIUSER — SQLcl connection Docker AI
The job, as the scheduler sees it.

select job_name, enabled, state, last_start_date, next_run_date
from   user_scheduler_jobs
where  job_name = 'COLDCHAIN_KAFKA_LOAD_JOB';

🟦 Run as: AIUSER — SQLcl connection Docker AI
Our own log: one row, 3 records, OK.

select run_ts, records_loaded, status
from   coldchain_load_log
order  by run_id desc;

🟦 Run as: AIUSER — SQLcl connection Docker AI
The most recent scheduler runs.

select status, actual_start_date, run_duration
from   user_scheduler_job_run_details
where  job_name = 'COLDCHAIN_KAFKA_LOAD_JOB'
order  by log_id desc
fetch  first 3 rows only;
DBMS_SCHEDULER
     |
     v
 coldchain_load_kafka
     |
     v
 DBMS_KAFKA.EXECUTE_LOAD_APP
     |
     v
 COLDCHAIN_READINGS

Two honest notes. First, OSAK keeps the Kafka offsets in Oracle metadata and advances them as it loads, so each call continues where the previous one stopped. Read Oracle's documentation for the exact transactional guarantees before you rely on them in production. Second, the log's run_ts is in database time, while the scheduler views show your session time zone, so the two look hours apart. They are the same moment.

Step 12: Prove it isn't just “above 8 is bad”

The database applies each shipment's own range. Send one acceptable reading and one that is too cold for the vaccines:

⬛ Run in: TERMINAL — PowerShell / Command Prompt on your machine
Start the producer.

docker exec -it kafka /opt/kafka/bin/kafka-console-producer.sh --bootstrap-server kafka:9092 --topic coldchain.readings

⬛ Run in: TERMINAL — PowerShell / Command Prompt on your machine
Paste, Enter, Ctrl+C.

{"sensor_id":"S-001","shipment_id":"SHP-1001","temp_c":7.2,"humidity_pct":62.0,"ts":"2026-10-03T07:20:00"}
{"sensor_id":"S-001","shipment_id":"SHP-1001","temp_c":1.5,"humidity_pct":62.0,"ts":"2026-10-03T07:21:00"}

🟦 Run as: AIUSER — SQLcl connection Docker AI
Wait a few seconds. Expect 1.5 = TOO_COLD and 7.2 = OK at the top.

select reading_ts, shipment_id, temp_c, min_temp_c, max_temp_c, status
from   coldchain_readings_v
where  shipment_id = 'SHP-1001'
order  by reading_ts desc, kafka_offset desc
fetch  first 5 rows only;

🟦 Run as: AIUSER — SQLcl connection Docker AI
Expect OK = 11, TOO_HOT = 6, TOO_COLD = 1 (18 readings).

select status, count(*) readings
from   coldchain_readings_v
group  by status
order  by status;

7.2°C → OK. 9.1°C → TOO_HOT. 1.5°C → TOO_COLD. Kafka carries the events; Oracle turns them into business information.

The honest limitation, and what comes next

PART 1:   Kafka  ------------------>  Oracle      (DBMS_KAFKA)

PART 2:   Oracle ------------------>  Kafka       (not DBMS_KAFKA)

DBMS_KAFKA is the ingestion side. It is not a Kafka producer API. When Oracle detects a breach and you want to publish an alert event back to Kafka, you need a different mechanism. That is Part 2.

Troubleshooting sheet

Kafka not reachable from Oracle

⬛ Run in: TERMINAL — PowerShell / Command Prompt on your machine
Are both containers up, and on the same network?

docker ps
docker network inspect kafka-net

⬛ Run in: TERMINAL — PowerShell / Command Prompt on your machine
Re-run the reachability test.

docker exec oracle26ai-apex261 bash -c "timeout 5 bash -c '</dev/tcp/kafka/9092' && echo REACHABLE || echo UNREACHABLE"

After restarting Docker Desktop, the cluster reconnected by itself in my test, but check the state anyway.

Cluster not connected

🟦 Run as: AIUSER — SQLcl connection Docker AI
Works as AIUSER once it holds OSAK_ADMIN_ROLE (or run it as SYS). Expect 0.

select DBMS_KAFKA_ADM.CHECK_CLUSTER('COLDCHAIN') as check_rc from dual;

select cluster_name, state, bootstrap_servers from dbms_kafka_clusters;

The load returns 0 records

🟦 Run as: AIUSER — SQLcl connection Docker AI
Is anything already loaded? If yes, produce a brand-new message and try again.

select count(*) from coldchain_readings;

Remember: never query the ORA$DKV_... views while troubleshooting.

The scheduler job isn't running

🟦 Run as: AIUSER — SQLcl connection Docker AI
State, next run, and the last five run outcomes.

select job_name, enabled, state, last_start_date, next_run_date
from   user_scheduler_jobs
where  job_name = 'COLDCHAIN_KAFKA_LOAD_JOB';

select status, actual_start_date, run_duration, additional_info
from   user_scheduler_job_run_details
where  job_name = 'COLDCHAIN_KAFKA_LOAD_JOB'
order  by log_id desc
fetch  first 5 rows only;

select * from coldchain_load_log where status = 'ERROR' order by run_id desc;

Pause and resume the stream

🟦 Run as: AIUSER — SQLcl connection Docker AI
Pause the stream, then resume it.

exec dbms_scheduler.disable('COLDCHAIN_KAFKA_LOAD_JOB');
exec dbms_scheduler.enable('COLDCHAIN_KAFKA_LOAD_JOB');

Reset: back to a clean slate for the next take

Run these in this order: AIUSER first (the load app must be dropped before the cluster is deregistered), then SYS, then the terminal.

Reset 1: AIUSER

🟦 Run as: AIUSER — SQLcl connection Docker AI
Drops the job, procedure, tables, view and load application.

begin
  begin dbms_scheduler.drop_job('COLDCHAIN_KAFKA_LOAD_JOB', force => true); exception when others then null; end;
  begin execute immediate 'drop procedure coldchain_load_kafka'; exception when others then null; end;
  begin execute immediate 'drop table coldchain_load_log purge'; exception when others then null; end;
  begin dbms_kafka.drop_load_app('COLDCHAIN', 'CCLOAD'); exception when others then null; end;
  begin execute immediate 'drop view coldchain_readings_v'; exception when others then null; end;
  begin execute immediate 'drop table coldchain_readings purge'; exception when others then null; end;
  begin execute immediate 'drop table coldchain_shipments purge'; exception when others then null; end;
end;
/

Reset 2: SYS

🟥 Run as: SYS — SQLcl connection Local Sys DBA
Clears the job's run history, deregisters the cluster, drops directories, revokes the grants.

begin
  dbms_scheduler.purge_log(log_history => 0, which_log => 'JOB_LOG',
                           job_name => 'AIUSER.COLDCHAIN_KAFKA_LOAD_JOB');
end;
/

begin
  dbms_kafka_adm.deregister_cluster('COLDCHAIN');
end;
/

drop directory OSAK_COLDCHAIN_ACCESS;
drop directory OSAK_COLDCHAIN_CONFIG;
revoke OSAK_ADMIN_ROLE from aiuser;
revoke create job from aiuser;

Reset 3: TERMINAL

⬛ Run in: TERMINAL — PowerShell / Command Prompt on your machine
Delete the topic (this throws away every test message).

docker exec kafka /opt/kafka/bin/kafka-topics.sh --bootstrap-server kafka:9092 --delete --topic coldchain.readings

⬛ Run in: TERMINAL — PowerShell / Command Prompt on your machine
Recreate it empty, so the demo counts start from zero again.

docker exec kafka /opt/kafka/bin/kafka-topics.sh --bootstrap-server kafka:9092 --create --if-not-exists --topic coldchain.readings --partitions 3 --replication-factor 1

⬛ Run in: TERMINAL — PowerShell / Command Prompt on your machine
Remove the OS folder inside the Oracle container (Step 4a recreates it).

docker exec oracle26ai-apex261 rm -rf /opt/oracle/osak

Then run the Step 0 checks again. Everything should read 0.

Build from scratch: the Docker pieces

Only needed if you don't already have the containers. Assumes your Oracle 26ai Free container is already running and named oracle26ai-apex261.

⬛ Run in: TERMINAL — PowerShell / Command Prompt on your machine
A user-defined network, so containers can find each other by name. Attach Oracle to it (it stays on the default network too, so its ports don't change).

docker network create kafka-net
docker network connect kafka-net oracle26ai-apex261

⬛ Run in: TERMINAL — PowerShell / Command Prompt on your machine
Pull the Kafka image first. The first download can take several minutes.

docker pull apache/kafka:3.9.0

Then start a single-node Kafka (KRaft mode, no ZooKeeper). The key detail is the two advertised listeners: kafka:9092 for containers on the network and localhost:29092 for your own machine.

⬛ Run in: TERMINAL — PowerShell / Command Prompt on your machine
One line, copy as is.

docker run -d --name kafka --network kafka-net --hostname kafka -p 29092:29092 -e KAFKA_NODE_ID=1 -e KAFKA_PROCESS_ROLES=broker,controller -e KAFKA_LISTENERS=INTERNAL://:9092,EXTERNAL://:29092,CONTROLLER://:9093 -e KAFKA_ADVERTISED_LISTENERS=INTERNAL://kafka:9092,EXTERNAL://localhost:29092 -e KAFKA_LISTENER_SECURITY_PROTOCOL_MAP=INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT,CONTROLLER:PLAINTEXT -e KAFKA_INTER_BROKER_LISTENER_NAME=INTERNAL -e KAFKA_CONTROLLER_LISTENER_NAMES=CONTROLLER -e KAFKA_CONTROLLER_QUORUM_VOTERS=1@kafka:9093 -e KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1 -e KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR=1 -e KAFKA_TRANSACTION_STATE_LOG_MIN_ISR=1 -e KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS=0 apache/kafka:3.9.0

Then continue with Step 0.

Gotchas collected along the way

  • ORA-62771: application names allow only letters, digits and $. No underscores.
  • ORA-62828: the load target table may contain only columns that exist in the OSAK view.
  • Package constants such as KAFKA_PROVIDER_APACHE can't be used inside a SQL select; pass the literal 'APACHE'.
  • A broker that advertises only localhost is unreachable from other containers. Advertise kafka:9092 as well.
  • Scheduler run history survives dropping and recreating a job, so counts can look inflated. Purge it as SYS (see Reset 2).
  • Don't query the generated ORA$DKV_... views directly.

No comments:

Powered by Blogger.