Skip to content
Projects
Groups
Snippets
Help
This project
Loading...
Sign in / Register
Toggle navigation
C
cdc-engine2
Overview
Overview
Details
Activity
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
0
Issues
0
List
Board
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
黄营
cdc-engine2
Commits
56eb47c3
Commit
56eb47c3
authored
Nov 12, 2024
by
y1sa
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
新增事件机制
parent
29f7a9ad
Hide whitespace changes
Inline
Side-by-side
Showing
2 changed files
with
12 additions
and
4 deletions
+12
-4
AbstractCdcEngine.java
...main/java/com/tbyf/cdcengine2/core/AbstractCdcEngine.java
+3
-3
ChangeListener.java
...rc/main/java/com/tbyf/cdcengine2/core/ChangeListener.java
+9
-1
No files found.
core/src/main/java/com/tbyf/cdcengine2/core/AbstractCdcEngine.java
View file @
56eb47c3
...
...
@@ -70,7 +70,7 @@ public abstract class AbstractCdcEngine<T extends AbstractCdcEngine<T>> {
ChangedRecord
record
=
recordQueue
.
take
();
try
{
for
(
ChangeListener
listener
:
changeListeners
)
{
if
(
listener
.
supports
(
record
))
{
if
(
listener
.
supports
(
record
.
schema
(),
record
.
table
()
))
{
listener
.
onChange
(
record
);
}
}
...
...
@@ -81,7 +81,7 @@ public abstract class AbstractCdcEngine<T extends AbstractCdcEngine<T>> {
for
(
ChangedRecord
record
:
buffer
)
{
try
{
for
(
ChangeListener
listener
:
changeListeners
)
{
if
(
listener
.
supports
(
record
))
{
if
(
listener
.
supports
(
record
.
schema
(),
record
.
table
()
))
{
listener
.
onChange
(
record
);
}
}
...
...
@@ -102,7 +102,7 @@ public abstract class AbstractCdcEngine<T extends AbstractCdcEngine<T>> {
for
(
ChangedRecord
record
:
buffer
)
{
try
{
for
(
ChangeListener
listener
:
changeListeners
)
{
if
(
listener
.
supports
(
record
))
{
if
(
listener
.
supports
(
record
.
schema
(),
record
.
table
()
))
{
listener
.
onChange
(
record
);
}
}
...
...
core/src/main/java/com/tbyf/cdcengine2/core/ChangeListener.java
View file @
56eb47c3
...
...
@@ -2,9 +2,17 @@ package com.tbyf.cdcengine2.core;
public
interface
ChangeListener
{
default
boolean
supports
(
ChangedRecord
record
)
{
default
boolean
supports
Schema
(
String
schema
)
{
return
true
;
}
default
boolean
supportsTable
(
String
table
)
{
return
true
;
}
default
boolean
supports
(
String
schema
,
String
table
)
{
return
supportsSchema
(
schema
)
&&
supportsTable
(
table
);
}
void
onChange
(
ChangedRecord
record
);
}
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment