You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
Copy file name to clipboardExpand all lines: data-prepper-plugins/rds-source/src/main/java/org/opensearch/dataprepper/plugins/source/rds/model/MessageType.java
Copy file name to clipboardExpand all lines: data-prepper-plugins/rds-source/src/main/java/org/opensearch/dataprepper/plugins/source/rds/stream/LogicalReplicationClient.java
+17-7
Original file line number
Diff line number
Diff line change
@@ -35,6 +35,7 @@ public class LogicalReplicationClient implements ReplicationLogClient {
Copy file name to clipboardExpand all lines: data-prepper-plugins/rds-source/src/main/java/org/opensearch/dataprepper/plugins/source/rds/stream/LogicalReplicationEventProcessor.java
+11-2
Original file line number
Diff line number
Diff line change
@@ -163,7 +163,16 @@ public void process(ByteBuffer msg) {
163
163
// If it's a RELATION, update table metadata map
164
164
// If it's INSERT/UPDATE/DELETE, prepare events
165
165
// If it's a COMMIT, convert all prepared events and send to buffer
Copy file name to clipboardExpand all lines: data-prepper-plugins/rds-source/src/main/java/org/opensearch/dataprepper/plugins/source/rds/stream/ReplicationLogClientFactory.java
+5-4
Original file line number
Diff line number
Diff line change
@@ -26,9 +26,9 @@
26
26
27
27
publicclassReplicationLogClientFactory {
28
28
29
-
privatefinalRdsSourceConfigsourceConfig;
30
29
privatefinalRdsClientrdsClient;
31
30
privatefinalDbMetadatadbMetadata;
31
+
privateRdsSourceConfigsourceConfig;
32
32
privateStringusername;
33
33
privateStringpassword;
34
34
privateSSLModesslMode = SSLMode.REQUIRED;
@@ -82,9 +82,10 @@ public void setSSLMode(SSLMode sslMode) {
Copy file name to clipboardExpand all lines: data-prepper-plugins/rds-source/src/main/java/org/opensearch/dataprepper/plugins/source/rds/stream/StreamWorker.java
+9-5
Original file line number
Diff line number
Diff line change
@@ -66,11 +66,15 @@ public void processStream(final StreamPartition streamPartition) {
Copy file name to clipboardExpand all lines: data-prepper-plugins/rds-source/src/main/java/org/opensearch/dataprepper/plugins/source/rds/stream/StreamWorkerTaskRefresher.java
+4-3
Original file line number
Diff line number
Diff line change
@@ -48,6 +48,7 @@ public class StreamWorkerTaskRefresher implements PluginConfigObserver<RdsSource
Copy file name to clipboardExpand all lines: data-prepper-plugins/rds-source/src/test/java/org/opensearch/dataprepper/plugins/source/rds/stream/LogicalReplicationClientTest.java
Copy file name to clipboardExpand all lines: data-prepper-plugins/rds-source/src/test/java/org/opensearch/dataprepper/plugins/source/rds/stream/LogicalReplicationEventProcessorTest.java
Copy file name to clipboardExpand all lines: data-prepper-plugins/rds-source/src/test/java/org/opensearch/dataprepper/plugins/source/rds/stream/StreamWorkerTest.java
0 commit comments