2020import com .navercorp .pinpoint .common .server .cluster .zookeeper .ZookeeperClient ;
2121import com .navercorp .pinpoint .common .server .cluster .zookeeper .ZookeeperEventWatcher ;
2222import com .navercorp .pinpoint .common .util .NetUtils ;
23- import com .navercorp .pinpoint .flink .config .FlinkConfiguration ;
23+ import com .navercorp .pinpoint .flink .config .FlinkProperties ;
2424import com .navercorp .pinpoint .rpc .util .ClassUtils ;
2525import com .navercorp .pinpoint .rpc .util .TimerFactory ;
2626import com .navercorp .pinpoint .web .cluster .zookeeper .PushZnodeJob ;
2727import com .navercorp .pinpoint .web .cluster .zookeeper .ZookeeperClusterDataManagerHelper ;
28-
2928import org .apache .curator .utils .ZKPaths ;
29+ import org .apache .logging .log4j .LogManager ;
30+ import org .apache .logging .log4j .Logger ;
3031import org .apache .zookeeper .WatchedEvent ;
3132import org .apache .zookeeper .Watcher .Event .EventType ;
3233import org .apache .zookeeper .Watcher .Event .KeeperState ;
3334import org .jboss .netty .util .HashedWheelTimer ;
3435import org .jboss .netty .util .Timeout ;
3536import org .jboss .netty .util .Timer ;
36- import org .apache .logging .log4j .Logger ;
37- import org .apache .logging .log4j .LogManager ;
3837
3938import javax .annotation .PostConstruct ;
4039import javax .annotation .PreDestroy ;
@@ -60,18 +59,18 @@ public class FlinkServerRegister implements ZookeeperEventWatcher {
6059
6160 private Timer timer ;
6261
63- public FlinkServerRegister (FlinkConfiguration flinkConfiguration ) {
64- Objects .requireNonNull (flinkConfiguration , "flinkConfiguration" );
65- this .clusterEnable = flinkConfiguration .isFlinkClusterEnable ();
66- this .connectAddress = flinkConfiguration .getFlinkClusterZookeeperAddress ();
67- this .sessionTimeout = flinkConfiguration .getFlinkClusterSessionTimeout ();
68- this .zookeeperPath = flinkConfiguration .getFlinkZNodePath ();
62+ public FlinkServerRegister (FlinkProperties flinkProperties ) {
63+ Objects .requireNonNull (flinkProperties , "flinkConfiguration" );
64+ this .clusterEnable = flinkProperties .isFlinkClusterEnable ();
65+ this .connectAddress = flinkProperties .getFlinkClusterZookeeperAddress ();
66+ this .sessionTimeout = flinkProperties .getFlinkClusterSessionTimeout ();
67+ this .zookeeperPath = flinkProperties .getFlinkZNodePath ();
6968
70- String zNodeName = getRepresentationLocalV4Ip () + ":" + flinkConfiguration .getFlinkClusterTcpPort ();
69+ String zNodeName = getRepresentationLocalV4Ip () + ":" + flinkProperties .getFlinkClusterTcpPort ();
7170 String zNodeFullPath = ZKPaths .makePath (zookeeperPath , zNodeName );
7271
7372 CreateNodeMessage createNodeMessage = new CreateNodeMessage (zNodeFullPath , new byte [0 ]);
74- int retryInterval = flinkConfiguration .getFlinkRetryInterval ();
73+ int retryInterval = flinkProperties .getFlinkRetryInterval ();
7574 this .pushFlinkNodeJob = new PushFlinkNodeJob (createNodeMessage , retryInterval );
7675 }
7776
0 commit comments