MCPcopy Create free account
hub / github.com/apache/kafka / start_migration_tool

Function start_migration_tool

system_test/utils/kafka_system_test_utils.py:1690–1752  ·  view source on GitHub ↗
(systemTestEnv, testcaseEnv, onlyThisEntityId=None)

Source from the content-addressed store, hash-verified

1688
1689
1690def start_migration_tool(systemTestEnv, testcaseEnv, onlyThisEntityId=None):
1691 clusterConfigList = systemTestEnv.clusterEntityConfigDictList
1692 migrationToolConfigList = system_test_utils.get_dict_from_list_of_dicts(clusterConfigList, "role", "migration_tool")
1693
1694 for migrationToolConfig in migrationToolConfigList:
1695
1696 entityId = migrationToolConfig["entity_id"]
1697
1698 if onlyThisEntityId is None or entityId == onlyThisEntityId:
1699
1700 host = migrationToolConfig["hostname"]
1701 jmxPort = migrationToolConfig["jmx_port"]
1702 role = migrationToolConfig["role"]
1703 kafkaHome = system_test_utils.get_data_by_lookup_keyval(clusterConfigList, "entity_id", entityId, "kafka_home")
1704 javaHome = system_test_utils.get_data_by_lookup_keyval(clusterConfigList, "entity_id", entityId, "java_home")
1705 jmxPort = system_test_utils.get_data_by_lookup_keyval(clusterConfigList, "entity_id", entityId, "jmx_port")
1706 kafkaRunClassBin = kafkaHome + "/bin/kafka-run-class.sh"
1707
1708 logger.info("starting kafka migration tool", extra=d)
1709 migrationToolLogPath = get_testcase_config_log_dir_pathname(testcaseEnv, "migration_tool", entityId, "default")
1710 migrationToolLogPathName = migrationToolLogPath + "/migration_tool.log"
1711 testcaseEnv.userDefinedEnvVarDict["migrationToolLogPathName"] = migrationToolLogPathName
1712
1713 testcaseConfigsList = testcaseEnv.testcaseConfigsList
1714 numProducers = system_test_utils.get_data_by_lookup_keyval(testcaseConfigsList, "entity_id", entityId, "num.producers")
1715 numStreams = system_test_utils.get_data_by_lookup_keyval(testcaseConfigsList, "entity_id", entityId, "num.streams")
1716 producerConfig = system_test_utils.get_data_by_lookup_keyval(testcaseConfigsList, "entity_id", entityId, "producer.config")
1717 consumerConfig = system_test_utils.get_data_by_lookup_keyval(testcaseConfigsList, "entity_id", entityId, "consumer.config")
1718 zkClientJar = system_test_utils.get_data_by_lookup_keyval(testcaseConfigsList, "entity_id", entityId, "zkclient.01.jar")
1719 kafka07Jar = system_test_utils.get_data_by_lookup_keyval(testcaseConfigsList, "entity_id", entityId, "kafka.07.jar")
1720 whiteList = system_test_utils.get_data_by_lookup_keyval(testcaseConfigsList, "entity_id", entityId, "whitelist")
1721 logFile = system_test_utils.get_data_by_lookup_keyval(testcaseConfigsList, "entity_id", entityId, "log_filename")
1722
1723 cmdList = ["ssh " + host,
1724 "'JAVA_HOME=" + javaHome,
1725 "JMX_PORT=" + jmxPort,
1726 kafkaRunClassBin + " kafka.tools.KafkaMigrationTool",
1727 "--whitelist=" + whiteList,
1728 "--num.producers=" + numProducers,
1729 "--num.streams=" + numStreams,
1730 "--producer.config=" + systemTestEnv.SYSTEM_TEST_BASE_DIR + "/" + producerConfig,
1731 "--consumer.config=" + systemTestEnv.SYSTEM_TEST_BASE_DIR + "/" + consumerConfig,
1732 "--zkclient.01.jar=" + systemTestEnv.SYSTEM_TEST_BASE_DIR + "/" + zkClientJar,
1733 "--kafka.07.jar=" + systemTestEnv.SYSTEM_TEST_BASE_DIR + "/" + kafka07Jar,
1734 " &> " + migrationToolLogPath + "/migrationTool.log",
1735 " & echo pid:$! > " + migrationToolLogPath + "/entity_" + entityId + "_pid'"]
1736
1737 cmdStr = " ".join(cmdList)
1738 logger.debug("executing command: [" + cmdStr + "]", extra=d)
1739 system_test_utils.async_sys_call(cmdStr)
1740 time.sleep(5)
1741
1742 pidCmdStr = "ssh " + host + " 'cat " + migrationToolLogPath + "/entity_" + entityId + "_pid' 2> /dev/null"
1743 logger.debug("executing command: [" + pidCmdStr + "]", extra=d)
1744 subproc = system_test_utils.sys_call_return_subproc(pidCmdStr)
1745
1746 # keep track of the remote entity pid in a dictionary
1747 for line in subproc.stdout.readlines():

Callers

nothing calls this directly

Calls 4

infoMethod · 0.80
joinMethod · 0.80
sleepMethod · 0.65

Tested by

no test coverage detected