Skip to content
Permalink
Browse files
Added some logging statements.
  • Loading branch information
nabarunnag committed Mar 3, 2020
1 parent fffb643 commit a713123163c10f4556896105c2a4b3e74eb6d4f2
Showing 2 changed files with 6 additions and 0 deletions.
@@ -56,6 +56,7 @@ public String version() {

@Override
public void start(Map<String, String> props) {
logger.info("Starting Apache Geode sink task");
try {
GeodeSinkConnectorConfig geodeConnectorConfig = new GeodeSinkConnectorConfig(props);
configure(geodeConnectorConfig);
@@ -142,6 +143,7 @@ private Region<Object, Object> createProxyRegion(String regionName) {

@Override
public void stop() {
logger.info("Stopping task");
geodeContext.close(false);
}

@@ -67,6 +67,7 @@ public String version() {

@Override
public void start(Map<String, String> props) {
logger.info("Starting Apache Geode source task");
try {
GeodeSourceConnectorConfig geodeConnectorConfig = new GeodeSourceConnectorConfig(props);
logger.debug("GeodeKafkaSourceTask id:" + geodeConnectorConfig.getTaskId() + " starting");
@@ -89,6 +90,7 @@ public void start(Map<String, String> props) {
boolean loadEntireRegion = geodeConnectorConfig.getLoadEntireRegion();
installOnGeode(geodeConnectorConfig, geodeContext, eventBufferSupplier, cqPrefix,
loadEntireRegion);
logger.info("Started Apache Geode source task");
} catch (Exception e) {
e.printStackTrace();
logger.error("Unable to start source task", e);
@@ -98,6 +100,7 @@ public void start(Map<String, String> props) {

@Override
public List<SourceRecord> poll() {
logger.trace("Polling for new data");
ArrayList<SourceRecord> records = new ArrayList<>(batchSize);
ArrayList<GeodeEvent> events = new ArrayList<>(batchSize);
if (eventBufferSupplier.get().drainTo(events, batchSize) > 0) {
@@ -117,6 +120,7 @@ public List<SourceRecord> poll() {

@Override
public void stop() {
logger.info("Stopping Apache Geode source task");
geodeContext.close(true);
}

0 comments on commit a713123

Please sign in to comment.