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
{{ message }}
This repository has been archived by the owner on Dec 14, 2022. It is now read-only.
val env = StreamExecutionEnvironment.getExecutionEnvironment
env.setParallelism(1)
env.enableCheckpointing(10000)
val prop = new Properties()
prop.put("topic", "persistent://ts/ns/mpc_test")
prop.put("pulsar.reader.subscriptionName", "ts.ns-profile")
val clientConf = new ClientConfigurationData()
clientConf.setAuthPluginClassName("org.apache.pulsar.client.impl.auth.AuthenticationToken")
clientConf.setAuthParams("eDkaliAfeGMsevPF0km6c1B741gxAWT2hgc")
clientConf.setServiceUrl("pulsar://127.0.0.1:6650")
val source1 = new FlinkPulsarSource[String]("http://127.0.0.1:8080", clientConf, new SimpleStringSchema(), prop)
source1.setStartFromSubscription(prop.getProperty("pulsar.reader.subscriptionName", "ts.ns-profile"))
val dd = env.addSource(source1)
dd.print("pulsar:").setParallelism(1)
env.execute(this.getClass.getName)
The text was updated successfully, but these errors were encountered:
val env = StreamExecutionEnvironment.getExecutionEnvironment
env.setParallelism(1)
env.enableCheckpointing(10000)
The text was updated successfully, but these errors were encountered: