How to use pyflink - 10 common examples

To help you get started, we’ve selected a few pyflink examples, based on popular ways it is used in public projects.

Secure your code as it's written. Use Snyk Code to scan source code in minutes - no build needed - and fix issues immediately.

github apache / flink / flink-python / pyflink / testing / source_sink_utils.py View on Github external
def __init__(self, field_names, field_types):
        TestTableSink._ensure_initialized()

        gateway = get_gateway()
        super(TestRetractSink, self).__init__(
            gateway.jvm.TestRetractSink(), field_names, field_types)
github apache / flink / flink-python / pyflink / testing / source_sink_utils.py View on Github external
def __init__(self, j_table_sink, field_names, field_types):
        gateway = get_gateway()
        j_field_names = utils.to_jarray(gateway.jvm.String, field_names)
        j_field_types = utils.to_jarray(gateway.jvm.TypeInformation,
                                        [_to_java_type(field_type) for field_type in field_types])
        j_table_sink = j_table_sink.configure(j_field_names, j_field_types)
        super(TestTableSink, self).__init__(j_table_sink)
github apache / flink / flink-python / pyflink / testing / test_case_utils.py View on Github external
def setUp(self):
        super(PyFlinkStreamTableTestCase, self).setUp()
        self.env = StreamExecutionEnvironment.get_execution_environment()
        self.env.set_parallelism(2)
        self.t_env = StreamTableEnvironment.create(self.env)
github apache / flink / flink-python / pyflink / testing / test_case_utils.py View on Github external
def setUp(self):
        super(PyFlinkBlinkStreamTableTestCase, self).setUp()
        self.env = StreamExecutionEnvironment.get_execution_environment()
        self.env.set_parallelism(2)
        self.t_env = StreamTableEnvironment.create(
            self.env, environment_settings=EnvironmentSettings.new_instance()
                .in_streaming_mode().use_blink_planner().build())
github apache / flink / flink-python / pyflink / testing / source_sink_utils.py View on Github external
def __init__(self, j_table_sink, field_names, field_types):
        gateway = get_gateway()
        j_field_names = utils.to_jarray(gateway.jvm.String, field_names)
        j_field_types = utils.to_jarray(gateway.jvm.TypeInformation,
                                        [_to_java_type(field_type) for field_type in field_types])
        j_table_sink = j_table_sink.configure(j_field_names, j_field_types)
        super(TestTableSink, self).__init__(j_table_sink)
github apache / flink / flink-python / pyflink / testing / test_case_utils.py View on Github external
def setUp(self):
        super(PyFlinkBatchTableTestCase, self).setUp()
        self.env = ExecutionEnvironment.get_execution_environment()
        self.env.set_parallelism(2)
        self.t_env = BatchTableEnvironment.create(self.env)
github apache / flink / flink-python / pyflink / testing / test_case_utils.py View on Github external
def setUp(self):
        super(PyFlinkBlinkBatchTableTestCase, self).setUp()
        self.t_env = BatchTableEnvironment.create(
            environment_settings=EnvironmentSettings.new_instance()
            .in_batch_mode().use_blink_planner().build())
        self.t_env._j_tenv.getPlanner().getExecEnv().setParallelism(2)
github apache / flink / flink-python / pyflink / testing / test_case_utils.py View on Github external
def setUp(self):
        super(PyFlinkBatchTableTestCase, self).setUp()
        self.env = ExecutionEnvironment.get_execution_environment()
        self.env.set_parallelism(2)
        self.t_env = BatchTableEnvironment.create(self.env)
github apache / flink / flink-python / pyflink / testing / test_case_utils.py View on Github external
def setUp(self):
        super(PyFlinkBlinkBatchTableTestCase, self).setUp()
        self.t_env = BatchTableEnvironment.create(
            environment_settings=EnvironmentSettings.new_instance()
            .in_batch_mode().use_blink_planner().build())
        self.t_env._j_tenv.getPlanner().getExecEnv().setParallelism(2)
github apache / flink / flink-python / pyflink / testing / test_case_utils.py View on Github external
def setUp(self):
        super(PyFlinkBlinkStreamTableTestCase, self).setUp()
        self.env = StreamExecutionEnvironment.get_execution_environment()
        self.env.set_parallelism(2)
        self.t_env = StreamTableEnvironment.create(
            self.env, environment_settings=EnvironmentSettings.new_instance()
                .in_streaming_mode().use_blink_planner().build())