writer.py 5.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141
  1. #!/usr/bin/env python
  2. # -*- coding: utf-8 -*-
  3. #
  4. # Copyright 2019 The FATE Authors. All Rights Reserved.
  5. #
  6. # Licensed under the Apache License, Version 2.0 (the "License");
  7. # you may not use this file except in compliance with the License.
  8. # You may obtain a copy of the License at
  9. #
  10. # http://www.apache.org/licenses/LICENSE-2.0
  11. #
  12. # Unless required by applicable law or agreed to in writing, software
  13. # distributed under the License is distributed on an "AS IS" BASIS,
  14. # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  15. # See the License for the specific language governing permissions and
  16. # limitations under the License.
  17. #
  18. from fate_arch.common import log
  19. from fate_arch.common.data_utils import default_output_fs_path
  20. from fate_arch.session import Session
  21. from fate_arch.storage import StorageEngine
  22. from fate_flow.components._base import (
  23. BaseParam,
  24. ComponentBase,
  25. ComponentInputProtocol,
  26. ComponentMeta,
  27. )
  28. from fate_flow.entity import Metric
  29. from fate_flow.external.data_storage import save_data_to_external_storage
  30. from fate_flow.manager.data_manager import DataTableTracker
  31. LOGGER = log.getLogger()
  32. writer_cpn_meta = ComponentMeta("Writer")
  33. @writer_cpn_meta.bind_param
  34. class WriterParam(BaseParam):
  35. def __init__(self,
  36. name=None,
  37. namespace=None,
  38. storage_engine=None,
  39. address=None,
  40. output_name=None,
  41. output_namespace=None):
  42. self.name = name
  43. self.namespace = namespace
  44. self.storage_engine = storage_engine
  45. self.address = address
  46. self.output_name = output_name
  47. self.output_namespace = output_namespace
  48. def check(self):
  49. return True
  50. @writer_cpn_meta.bind_runner.on_guest.on_host.on_local
  51. class Writer(ComponentBase):
  52. def __init__(self):
  53. super(Writer, self).__init__()
  54. self.parameters = None
  55. self.job_parameters = None
  56. def _run(self, cpn_input: ComponentInputProtocol):
  57. self.parameters = cpn_input.parameters
  58. if self.parameters.get("namespace") and self.parameters.get("name"):
  59. namespace = self.parameters.get("namespace")
  60. name = self.parameters.get("name")
  61. elif cpn_input.flow_feeded_parameters.get("table_info"):
  62. namespace = cpn_input.flow_feeded_parameters.get("table_info")[0].get("namespace")
  63. name = cpn_input.flow_feeded_parameters.get("table_info")[0].get("name")
  64. else:
  65. raise Exception("no found name or namespace in input parameters")
  66. LOGGER.info(f"writer parameters:{self.parameters}")
  67. src_table = Session.get_global().get_table(name=name, namespace=namespace)
  68. output_name = self.parameters.get("output_name")
  69. output_namespace = self.parameters.get("output_namespace")
  70. engine = self.parameters.get("storage_engine")
  71. address_dict = self.parameters.get("address")
  72. if output_name and output_namespace:
  73. table_meta = src_table.meta.to_dict()
  74. address_dict = src_table.meta.get_address().__dict__
  75. engine = src_table.meta.get_engine()
  76. table_meta.update({
  77. "name": output_name,
  78. "namespace": output_namespace,
  79. "address": self._create_save_address(engine, address_dict, output_name, output_namespace),
  80. })
  81. src_table.save_as(**table_meta)
  82. table = src_table.save_as(**table_meta)
  83. table.meta.update_metas(**table_meta)
  84. # output table track
  85. DataTableTracker.create_table_tracker(
  86. name,
  87. namespace,
  88. entity_info={
  89. "have_parent": True,
  90. "parent_table_namespace": namespace,
  91. "parent_table_name": name,
  92. "job_id": self.tracker.job_id,
  93. }
  94. )
  95. elif engine and address_dict:
  96. save_data_to_external_storage(engine, address_dict, src_table)
  97. LOGGER.info("save success")
  98. self.tracker.log_output_data_info(
  99. data_name="writer",
  100. table_namespace=output_namespace,
  101. table_name=output_name,
  102. )
  103. self.tracker.log_metric_data(
  104. metric_namespace="writer",
  105. metric_name="writer",
  106. metrics=[Metric("count", src_table.meta.get_count()),
  107. Metric("storage_engine", engine)]
  108. )
  109. @staticmethod
  110. def _create_save_address(engine, address_dict, name, namespace):
  111. if engine == StorageEngine.EGGROLL:
  112. address_dict.update({"name": name, "namespace": namespace})
  113. elif engine == StorageEngine.STANDALONE:
  114. address_dict.update({"name": name, "namespace": namespace})
  115. elif engine == StorageEngine.HIVE:
  116. address_dict.update({"database": namespace, "name": f"{name}"})
  117. elif engine == StorageEngine.HDFS:
  118. address_dict.update({"path": default_output_fs_path(name=name,
  119. namespace=namespace,
  120. prefix=address_dict.get("path_prefix"))})
  121. elif engine == StorageEngine.LOCALFS:
  122. address_dict.update({"path": default_output_fs_path(name=name, namespace=namespace,
  123. storage_engine=StorageEngine.LOCALFS)})
  124. else:
  125. raise RuntimeError(f"{engine} storage is not supported")
  126. return address_dict