diff --git a/README.md b/README.md index 599f8f149e..4ee91d03bf 100644 --- a/README.md +++ b/README.md @@ -4,11 +4,10 @@ * **FlinkX是在是袋鼠云内部广泛使用的基于flink的分布式离线数据同步框架,实现了多种异构数据源之间高效的数据迁移。** - 不同的数据源头被抽象成不同的Reader插件,不同的数据目标被抽象成不同的Writer插件。理论上,FlinkX框架可以支持任意数据源类型的数据同步工作。作为一套生态系统,每接入一套新数据源该新加入的数据源即可实现和现有的数据源互通。
+
+
+ * http://www.apache.org/licenses/LICENSE-2.0 + *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package com.dtstack.flinkx.metrics;
+
+import com.dtstack.flinkx.constants.Metrics;
+import org.apache.flink.api.common.accumulators.LongCounter;
+import org.apache.flink.api.common.functions.RuntimeContext;
+import org.apache.flink.metrics.MetricGroup;
+import org.apache.flink.runtime.metrics.MetricRegistryImpl;
+import org.apache.flink.runtime.metrics.groups.AbstractMetricGroup;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import scala.concurrent.duration.FiniteDuration;
+
+import java.lang.reflect.Field;
+import java.util.concurrent.RunnableScheduledFuture;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.ScheduledThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+
+/**
+ * company: www.dtstack.com
+ *
+ * @author: toutian
+ * create: 2019/3/18
+ */
+public class InputMetric {
+ protected final Logger LOG = LoggerFactory.getLogger(getClass());
+
+ private RuntimeContext runtimeContext;
+
+ private final static Long DEFAULT_PERIOD_MILLISECONDS = 10000L;
+
+ private Long delayPeriodMill = 12000L;
+
+ public InputMetric(RuntimeContext runtimeContext, LongCounter numRead) {
+ this.runtimeContext = runtimeContext;
+
+ final MetricGroup flinkxInput = getRuntimeContext().getMetricGroup().addGroup(Metrics.METRIC_GROUP_KEY_FLINKX, Metrics.METRIC_GROUP_VALUE_INPUT);
+
+ flinkxInput.gauge(Metrics.NUM_READS, new SimpleAccumulatorGauge