Skip to content

Commit 5032c36

Browse files
authored
Merge pull request #16 from oracle-devrel/updateversions
Updateversions
2 parents dfacb1c + c4d4c73 commit 5032c36

24 files changed

Lines changed: 773 additions & 83 deletions

File tree

IoTDBJDBC/.factorypath

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,13 @@
11
<factorypath>
22
<factorypathentry kind="VARJAR" id="M2_REPO/org/projectlombok/lombok/1.18.42/lombok-1.18.42.jar" enabled="true" runInBatchMode="false"/>
3-
<factorypathentry kind="VARJAR" id="M2_REPO/io/micronaut/micronaut-inject-java/4.10.18/micronaut-inject-java-4.10.18.jar" enabled="true" runInBatchMode="false"/>
3+
<factorypathentry kind="VARJAR" id="M2_REPO/io/micronaut/micronaut-inject-java/4.10.23/micronaut-inject-java-4.10.23.jar" enabled="true" runInBatchMode="false"/>
44
<factorypathentry kind="VARJAR" id="M2_REPO/org/slf4j/slf4j-api/2.0.17/slf4j-api-2.0.17.jar" enabled="true" runInBatchMode="false"/>
5-
<factorypathentry kind="VARJAR" id="M2_REPO/io/micronaut/micronaut-core-processor/4.10.18/micronaut-core-processor-4.10.18.jar" enabled="true" runInBatchMode="false"/>
6-
<factorypathentry kind="VARJAR" id="M2_REPO/io/micronaut/micronaut-inject/4.10.18/micronaut-inject-4.10.18.jar" enabled="true" runInBatchMode="false"/>
5+
<factorypathentry kind="VARJAR" id="M2_REPO/io/micronaut/micronaut-core-processor/4.10.23/micronaut-core-processor-4.10.23.jar" enabled="true" runInBatchMode="false"/>
6+
<factorypathentry kind="VARJAR" id="M2_REPO/io/micronaut/micronaut-inject/4.10.23/micronaut-inject-4.10.23.jar" enabled="true" runInBatchMode="false"/>
77
<factorypathentry kind="VARJAR" id="M2_REPO/jakarta/inject/jakarta.inject-api/2.0.1/jakarta.inject-api-2.0.1.jar" enabled="true" runInBatchMode="false"/>
88
<factorypathentry kind="VARJAR" id="M2_REPO/jakarta/annotation/jakarta.annotation-api/2.1.1/jakarta.annotation-api-2.1.1.jar" enabled="true" runInBatchMode="false"/>
9-
<factorypathentry kind="VARJAR" id="M2_REPO/io/micronaut/micronaut-core/4.10.18/micronaut-core-4.10.18.jar" enabled="true" runInBatchMode="false"/>
10-
<factorypathentry kind="VARJAR" id="M2_REPO/io/micronaut/micronaut-aop/4.10.18/micronaut-aop-4.10.18.jar" enabled="true" runInBatchMode="false"/>
9+
<factorypathentry kind="VARJAR" id="M2_REPO/io/micronaut/micronaut-core/4.10.23/micronaut-core-4.10.23.jar" enabled="true" runInBatchMode="false"/>
10+
<factorypathentry kind="VARJAR" id="M2_REPO/io/micronaut/micronaut-aop/4.10.23/micronaut-aop-4.10.23.jar" enabled="true" runInBatchMode="false"/>
1111
<factorypathentry kind="VARJAR" id="M2_REPO/com/github/javaparser/javaparser-symbol-solver-core/3.27.0/javaparser-symbol-solver-core-3.27.0.jar" enabled="true" runInBatchMode="false"/>
1212
<factorypathentry kind="VARJAR" id="M2_REPO/com/github/javaparser/javaparser-core/3.27.0/javaparser-core-3.27.0.jar" enabled="true" runInBatchMode="false"/>
1313
<factorypathentry kind="VARJAR" id="M2_REPO/org/checkerframework/checker-qual/3.49.4/checker-qual-3.49.4.jar" enabled="true" runInBatchMode="false"/>
@@ -17,10 +17,10 @@
1717
<factorypathentry kind="VARJAR" id="M2_REPO/org/ow2/asm/asm-tree/9.8/asm-tree-9.8.jar" enabled="true" runInBatchMode="false"/>
1818
<factorypathentry kind="VARJAR" id="M2_REPO/org/ow2/asm/asm-util/9.8/asm-util-9.8.jar" enabled="true" runInBatchMode="false"/>
1919
<factorypathentry kind="VARJAR" id="M2_REPO/org/ow2/asm/asm-analysis/9.8/asm-analysis-9.8.jar" enabled="true" runInBatchMode="false"/>
20-
<factorypathentry kind="VARJAR" id="M2_REPO/io/micronaut/micronaut-core-reactive/4.10.18/micronaut-core-reactive-4.10.18.jar" enabled="true" runInBatchMode="false"/>
20+
<factorypathentry kind="VARJAR" id="M2_REPO/io/micronaut/micronaut-core-reactive/4.10.23/micronaut-core-reactive-4.10.23.jar" enabled="true" runInBatchMode="false"/>
2121
<factorypathentry kind="VARJAR" id="M2_REPO/org/reactivestreams/reactive-streams/1.0.4/reactive-streams-1.0.4.jar" enabled="true" runInBatchMode="false"/>
2222
<factorypathentry kind="VARJAR" id="M2_REPO/org/ow2/asm/asm/9.8/asm-9.8.jar" enabled="true" runInBatchMode="false"/>
23-
<factorypathentry kind="VARJAR" id="M2_REPO/io/micronaut/micronaut-graal/4.10.18/micronaut-graal-4.10.18.jar" enabled="true" runInBatchMode="false"/>
23+
<factorypathentry kind="VARJAR" id="M2_REPO/io/micronaut/micronaut-graal/4.10.23/micronaut-graal-4.10.23.jar" enabled="true" runInBatchMode="false"/>
2424
<factorypathentry kind="VARJAR" id="M2_REPO/io/micronaut/serde/micronaut-serde-processor/2.16.2/micronaut-serde-processor-2.16.2.jar" enabled="true" runInBatchMode="false"/>
2525
<factorypathentry kind="VARJAR" id="M2_REPO/io/micronaut/serde/micronaut-serde-api/2.16.2/micronaut-serde-api-2.16.2.jar" enabled="true" runInBatchMode="false"/>
2626
<factorypathentry kind="VARJAR" id="M2_REPO/io/micronaut/micronaut-context/4.9.11/micronaut-context-4.9.11.jar" enabled="true" runInBatchMode="false"/>

IoTDBJDBC/README.md

Lines changed: 90 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -3,8 +3,97 @@ This is built using maven and the micronaut libraries. see below for details on
33

44
YOU MUST be running this on a compute resource on a private VCN that is connected to the IOT database, see the IOT documentation at https://docs.oracle.com/en-us/iaas/Content/internet-of-things/connect-database.htm for details on how to set this up.
55

6+
## Dynamic source and handler selection
67

7-
]## Micronaut 4.10.10 Documentation
8+
The application uses Micronaut dependency injection to decide which IoT data sources and output handlers are active at runtime. The code does not use a central switch statement. Instead, each optional reader, filter, processor, and output is a Micronaut `@Singleton` bean guarded with `@Requires` annotations. If the required configuration properties are present, and any `.enabled` property is set to `true`, Micronaut creates that bean. If the properties are missing or disabled, the bean is not created and is not injected anywhere else.
9+
10+
### Selecting input sources
11+
12+
All database and AQ readers implement `IoTDBClient`. `JDBCRunner` asks Micronaut for `List<IoTDBClient>`, so only the readers whose `@Requires` conditions pass are injected. The runner sorts the injected clients by each client's configured `getOrder()` value, calls `configureDBClient()`, then starts them with `startDBProcessing()`. This means one configuration file can enable one source, several sources, or none.
13+
14+
The available input clients are:
15+
16+
| Client | Data source | Main enable property | Order property |
17+
| --- | --- | --- | --- |
18+
| `IoTJDBCConnectionTestReader` | Reads sample rows from the `raw_data` table over JDBC. | `iotdatacache.jdbc.doconnectiontestread.enabled=true` | `iotdatacache.jdbc.doconnectiontestread.order` |
19+
| `IoTAQNormalizedDataBatchReader` | Dequeues batches from the normalized data AQ queue and passes each message to the normalized data handler chain. | `iotdatacache.aq.normalizeddata.batchreader.enabled=true` | `iotdatacache.aq.normalizeddata.batchreader.order` |
20+
| `IoTAQNormalizedDataIndividualReader` | Dequeues normalized AQ messages one at a time and passes them to the normalized data handler chain. | `iotdatacache.aq.normalizeddata.individualreader.enabled=true` | `iotdatacache.aq.normalizeddata.individualreader.order` |
21+
| `IoTAQRawDataIndividualReader` | Dequeues raw data AQ messages one at a time and passes them to the raw data handler chain. | `iotdatacache.aq.rawdata.individualreader.enabled=true` | `iotdatacache.aq.rawdata.individualreader.order` |
22+
| `IoTAQNormalizedDataListener` | Registers an AQ notification listener for normalized data. | `iotdatacache.aq.listener.enabled=true` | `iotdatacache.aq.listener.order` |
23+
24+
Readers can also have source-specific settings. For example, the batch reader uses `iotdatacache.aq.normalizeddata.batchreader.batchsize`, and the AQ readers use read timeout and subscriber-name properties. See `config/config.properties` for examples of the available settings.
25+
26+
The normalized batch, normalized individual, and raw individual AQ readers feed messages into the handler services. The notification listener currently dequeues and logs normalized messages through its core output method.
27+
28+
### AQ versus JDBC
29+
30+
The JDBC reader and the AQ readers both use database connections, but they use different database interaction models. The JDBC reader runs SQL directly against database tables, for example reading recent rows from `raw_data`. It is useful for direct queries, connection checks, diagnostics, and any case where the application wants to decide exactly which table rows to fetch.
31+
32+
AQ, or Advanced Queuing, treats incoming IoT data as queued messages rather than table rows to query. The AQ readers subscribe to the IoT queue, dequeue messages from it, convert the queue payload into `RawData` or `NormalizedData`, and then pass the converted object to the appropriate message handler service.
33+
34+
The current AQ readers can run as polling loops. For example, the individual and batch readers repeatedly call `dequeue` with a configured wait timeout. If no message is available before the timeout, the loop simply waits again. When a message is retrieved, the reader immediately hands it to `RawDataMessageHandlerService` or `NormalizedDataMessageHandlerService`. Because downstream code receives each retrieved message through the handler chain as soon as the polling loop obtains it, the rest of the application can treat the flow as listener-like even though the reader itself is implemented as a polling loop.
35+
36+
The AQ `dequeue` call is blocking up to the configured wait timeout. For a single-message dequeue, the call returns as soon as one message is available, or it times out if no message arrives. For a batch dequeue, the call requests up to the configured batch size and returns when the requested number of entries has been retrieved or when the wait timeout expires. That means a batch call may return a full batch, a partial batch, or no data if the queue stays empty until timeout.
37+
38+
Requesting one entry at a time keeps per-message latency low and makes error handling simple, because each dequeue result maps directly to one handler-chain invocation. It also commits progress frequently. The tradeoff is that it performs more database round trips when the queue is busy, so throughput can be lower.
39+
40+
Requesting multiple entries in one dequeue can improve throughput by reducing database round trips and allowing a batch reader to drain bursts of queued messages more efficiently. The tradeoff is that an early message in the batch may wait until the batch fills or the timeout fires, and processing/commit behavior is grouped around the batch read. Batch reads are usually better for sustained or bursty load; single-message reads are usually better when immediate processing and simpler operational behavior matter more than maximum throughput.
41+
42+
This gives the application two useful modes. Direct JDBC access is table/query oriented. AQ access is message/stream oriented and is better suited to continuous processing, filtering, transformation, and output handling.
43+
44+
### Selecting filters, processors, and outputs
45+
46+
Raw and normalized messages are handled by separate chains:
47+
48+
- `RawDataMessageHandlerService` receives a Micronaut-injected `List<RawDataMessageHandler>`.
49+
- `NormalizedDataMessageHandlerService` receives a Micronaut-injected `List<NormalizedDataMessageHandler>`.
50+
51+
As with input clients, only handlers whose `@Requires` properties match the runtime configuration are created and injected. Each handler provides an order value from configuration, and the service sorts the injected handlers before processing messages. A handler returns an array of messages. Returning one or more messages passes those messages to the next handler in the chain; returning an empty array stops that branch. This allows the same mechanism to support filters, test processors, and final output handlers.
52+
53+
The raw data handler chain can include filters such as:
54+
55+
- `messagehandler.filter.rawdata.contenttype.*`
56+
- `messagehandler.filter.rawdata.endpointfilter.*`
57+
- `messagehandler.filter.rawdata.devicemodelfilter.*`
58+
59+
It can also include outputs such as:
60+
61+
- text output using `messagehandler.output.rawdata.textoutput.*`
62+
- HTTP output enabled with `messagehandler.output.rawdata.httpclient.enabled`; the current code reads its order and type from `messagehandler.output.rawdata.httpclient.enabled.order` and `messagehandler.output.rawdata.httpclient.enabled.type`
63+
- NoSQL output settings under `messagehandler.output.rawdata.nosql.*`; the current code enables this bean with `messagehandler.filter.rawdata.nosql.enabled`
64+
65+
The normalized data handler chain can include filters and test processors such as:
66+
67+
- `messagehandler.filter.normalizeddata.contentpathfilter.*`
68+
- `messagehandler.filter.normalizeddata.devicemodelfilter.*`
69+
- `messagehandler.filter.normalizeddata.randomfilter.*`
70+
- `messagehandler.processor.normalizeddata.duplicator.*`
71+
72+
It can also include the diagnostic text output:
73+
74+
- `messagehandler.output.normalizeddata.textoutput.*`
75+
76+
The usual pattern is to set a handler's `.enabled` property to `true` and provide its `.order` property. Lower order values run earlier. Filters normally sit before outputs, and outputs can either pass the message through to later handlers or terminate that branch by returning no messages.
77+
78+
Example:
79+
80+
```properties
81+
iotdatacache.aq.rawdata.individualreader.enabled=true
82+
iotdatacache.aq.rawdata.individualreader.order=10
83+
iotdatacache.aq.rawdata.individualreader.readtimeout=10
84+
85+
messagehandler.filter.rawdata.contenttype.enabled=true
86+
messagehandler.filter.rawdata.contenttype.order=10
87+
messagehandler.filter.rawdata.contenttype.type=application/json
88+
89+
messagehandler.output.rawdata.textoutput.enabled=true
90+
messagehandler.output.rawdata.textoutput.order=20
91+
messagehandler.output.rawdata.textoutput.passthrough=false
92+
```
93+
94+
In this example, Micronaut creates the raw AQ reader, the content-type filter, and the text output handler. The reader receives raw AQ messages, the raw handler service runs the content-type filter first, and matching messages are then passed to the text output handler.
95+
96+
## Micronaut 4.10.10 Documentation
897

998
- [User Guide](https://docs.micronaut.io/4.10.10/guide/index.html)
1099
- [API Reference](https://docs.micronaut.io/4.10.10/api/index.html)
@@ -56,5 +145,3 @@ YOU MUST be running this on a compute resource on a private VCN that is connecte
56145

57146

58147
- [Micronaut AOT documentation](https://micronaut-projects.github.io/micronaut-aot/latest/guide/)
59-
60-

IoTDBJDBC/pom.xml

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -46,20 +46,20 @@ SOFTWARE. -->
4646
<parent>
4747
<groupId>io.micronaut.platform</groupId>
4848
<artifactId>micronaut-parent</artifactId>
49-
<version>4.10.10</version>
49+
<version>4.10.14</version>
5050
</parent>
5151

5252
<properties>
5353
<packaging>jar</packaging>
5454
<jdk.version>21</jdk.version>
5555
<release.version>17</release.version>
5656
<version.ocisdk>3.74.2</version.ocisdk>
57-
<version.lombok>1.18.44</version.lombok>
57+
<version.lombok>1.18.46</version.lombok>
5858
<version.slf4j>2.0.7</version.slf4j>
5959
<!-- Oracle JDBC 23ai line (set to the latest approved version) -->
6060
<version.ojdbc>23.8.0.25.04</version.ojdbc>
6161
<version.nosql>3.85.0</version.nosql>
62-
<micronaut.version>4.10.21</micronaut.version>
62+
<!--<micronaut.version>4.10.25</micronaut.version>-->
6363
<micronaut.runtime>netty</micronaut.runtime>
6464
<micronaut.aot.enabled>false</micronaut.aot.enabled>
6565
<micronaut.aot.packageName>

IoTDBJDBC/src/main/java/com/oracle/demo/timg/iot/iotdbjdbc/messagehandler/filters/DeviceModelMessageFilterCore.java

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -181,15 +181,16 @@ public boolean doesIoTDataCoreMatchModel(IoTDataCore input) throws Exception {
181181
}
182182
log.fine(() -> "instance is unknown retrieving its model, " + instanceId);
183183
// we don't know about it, using the device ID query the DB to get the model id
184-
String instanceModelId;
184+
String instanceModelIdTemp = null;
185185
try {
186-
instanceModelId = getModelIdFromInstanceId(instanceId);
186+
instanceModelIdTemp = getModelIdFromInstanceId(instanceId);
187187
} catch (SQLException e) {
188188
log.warning("SQLException locating instances model id for model " + instanceId + ", "
189189
+ e.getLocalizedMessage());
190-
instanceModelId = null;
191190
}
192-
// can't use a lambda here as instanceModelId is
191+
// can't use a lambda for the debugging unless we do this as instanceModelId
192+
// must be final
193+
String instanceModelId = instanceModelIdTemp;
193194
log.fine("instance has model id, " + instanceModelId);
194195
if (instanceModelId == null) {
195196
// no model id found, this I guess is possible for an instance that is not

IoTDBJDBC/src/main/java/com/oracle/demo/timg/iot/iotdbjdbc/messagehandler/outputs/http/IoTOutputHttiClientRequestFilter.java renamed to IoTDBJDBC/src/main/java/com/oracle/demo/timg/iot/iotdbjdbc/messagehandler/outputs/http/IoTOutputHttpClientRequestFilter.java

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -49,12 +49,12 @@ Software and the Larger Work(s), and to sublicense the foregoing rights on
4949

5050
@ClientFilter(patterns = { "${messagehandler.output.iotoutputhttpclient:/api/v1/iotdata}/**"})
5151
@Log
52-
public class IoTOutputHttiClientRequestFilter {
52+
public class IoTOutputHttpClientRequestFilter {
5353
private final String username;
5454
private final String password;
5555

5656
@Inject
57-
public IoTOutputHttiClientRequestFilter(IoTOutputHttpClientSettings clientSettings) {
57+
public IoTOutputHttpClientRequestFilter(IoTOutputHttpClientSettings clientSettings) {
5858
this.username = clientSettings.getUsername();
5959
this.password = new String(Base64.getDecoder().decode(clientSettings.getPassword()));
6060
}
@@ -67,6 +67,6 @@ public void doFilter(MutableHttpRequest<?> request) {
6767

6868
@EventListener
6969
public void onStartup(StartupEvent event) {
70-
log.info("Startup event received for IoTOutputHttiClientRequestFilter username=" + this.username);
70+
log.info("Startup event received for IoTOutputHttpClientRequestFilter username=" + this.username);
7171
}
7272
}

IoTDBJDBC/src/main/java/com/oracle/demo/timg/iot/iotdbjdbc/messagehandler/outputs/http/IoTOutputHttpClientSettings.java

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -43,10 +43,8 @@ Software and the Larger Work(s), and to sublicense the foregoing rights on
4343
@Data
4444
public class IoTOutputHttpClientSettings {
4545
public static final String PREFIX = "messagehandler.output.iotoutputhttpclient";
46-
// @Property(name = IoTOutputHttpClientSettings.PREFIX + ".username",
47-
// defaultValue = "")
46+
// as we are operating as a configuration then the fields are set based on the
47+
// config tree
4848
private String username;
49-
// @Property(name = IoTOutputHttpClientSettings.PREFIX + ".password",
50-
// defaultValue = "")
5149
private String password;
5250
}

0 commit comments

Comments
 (0)